From e05f8f86a104ab72ce596a04a6fc941b80f1856c Mon Sep 17 00:00:00 2001 From: Dmitry Batkovich Date: Fri, 17 Jun 2022 17:36:50 +0400 Subject: [PATCH] add lock-free persistent fs records storage implementation (disabled by default) GitOrigin-RevId: 2eb6dfd404cf4e994bec01f427d357b4c235534d --- .../openapi/vfs/newvfs/ManagingFS.java | 2 +- .../com/intellij/util/containers/Unsafe.java | 2 +- .../com/intellij/util/io/ByteBufferUtil.java | 9 + .../util/io/ResizeableMappedFile.java | 8 + .../intellij/util/io/StorageLockContext.java | 21 +- .../vfs/PlatformVirtualFileManager.java | 2 +- .../vfs/newvfs/persistent/FSRecords.java | 248 ++++++++--- .../persistent/PersistentFSConnection.java | 21 +- .../persistent/PersistentFSConnector.java | 19 +- .../persistent/PersistentFSHeaders.java | 3 +- .../newvfs/persistent/PersistentFSImpl.java | 2 +- .../PersistentFSLockFreeRecordsStorage.kt | 389 ++++++++++++++++++ .../PersistentFSRecordAccessor.java | 15 +- .../PersistentFSRecordsStorage.java | 334 ++------------- ...ersistentFSSynchronizedRecordsStorage.java | 368 +++++++++++++++++ 15 files changed, 1053 insertions(+), 390 deletions(-) create mode 100644 platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSLockFreeRecordsStorage.kt create mode 100644 platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSSynchronizedRecordsStorage.java diff --git a/platform/analysis-api/src/com/intellij/openapi/vfs/newvfs/ManagingFS.java b/platform/analysis-api/src/com/intellij/openapi/vfs/newvfs/ManagingFS.java index 0151eb80f06a..2187df3c7a86 100644 --- a/platform/analysis-api/src/com/intellij/openapi/vfs/newvfs/ManagingFS.java +++ b/platform/analysis-api/src/com/intellij/openapi/vfs/newvfs/ManagingFS.java @@ -48,7 +48,7 @@ public abstract class ManagingFS implements FileSystemInterface { /** * @return a number that's incremented every time modification count for some file is advanced, @see {@link #getModificationCount(VirtualFile)}. - * This number is persisted between IDE sessions and so it'll always increase. This method invocation means disk access, so it's not terribly cheap. + * This number is persisted between IDE sessions and so it'll always increase. */ @TestOnly public abstract int getFilesystemModificationCount(); diff --git a/platform/util/src/com/intellij/util/containers/Unsafe.java b/platform/util/src/com/intellij/util/containers/Unsafe.java index b6bf86770f8a..f5cf95cca86e 100644 --- a/platform/util/src/com/intellij/util/containers/Unsafe.java +++ b/platform/util/src/com/intellij/util/containers/Unsafe.java @@ -61,7 +61,7 @@ class Unsafe { .bindTo(unsafe); } - static boolean compareAndSwapInt(@NotNull Object object, long offset, int expected, int value) { + public static boolean compareAndSwapInt(Object object, long offset, int expected, int value) { try { return (boolean)compareAndSwapInt.invokeExact(object, offset, expected, value); } diff --git a/platform/util/src/com/intellij/util/io/ByteBufferUtil.java b/platform/util/src/com/intellij/util/io/ByteBufferUtil.java index fa308fce8d15..0394b2b354fd 100644 --- a/platform/util/src/com/intellij/util/io/ByteBufferUtil.java +++ b/platform/util/src/com/intellij/util/io/ByteBufferUtil.java @@ -108,6 +108,15 @@ public final class ByteBufferUtil { buf.get(dst, dstIndex, length); } + public static long getAddress(@NotNull ByteBuffer src) { + try { + return (long)address.invoke(src); + } + catch (Throwable e) { + throw new RuntimeException(e); + } + } + @NotNull private static Logger getLogger() { return Logger.getInstance(ByteBufferUtil.class); diff --git a/platform/util/src/com/intellij/util/io/ResizeableMappedFile.java b/platform/util/src/com/intellij/util/io/ResizeableMappedFile.java index c3f144d8198a..b2a9087dd603 100644 --- a/platform/util/src/com/intellij/util/io/ResizeableMappedFile.java +++ b/platform/util/src/com/intellij/util/io/ResizeableMappedFile.java @@ -248,6 +248,14 @@ public class ResizeableMappedFile implements Forceable, Closeable { myStorage.putBuffer(index, buffer); } + public void setLogicalSize(long logicalSize) { + myLogicalSize = logicalSize; + } + + public long getLogicalSize() { + return myLogicalSize; + } + public void close() throws IOException { List exceptions = new SmartList<>(); ContainerUtil.addIfNotNull(exceptions, ExceptionUtil.runAndCatch(() -> { diff --git a/platform/util/src/com/intellij/util/io/StorageLockContext.java b/platform/util/src/com/intellij/util/io/StorageLockContext.java index 1e16b2d10fd2..af8fe49524d9 100644 --- a/platform/util/src/com/intellij/util/io/StorageLockContext.java +++ b/platform/util/src/com/intellij/util/io/StorageLockContext.java @@ -19,27 +19,36 @@ public final class StorageLockContext { private final FilePageCache myFilePageCache; private final boolean myUseReadWriteLock; private final boolean myCacheChannels; + private final boolean myDisableAssertions; private StorageLockContext(@NotNull FilePageCache filePageCache, boolean useReadWriteLock, - boolean cacheChannels) { + boolean cacheChannels, + boolean disableAssertions) { myLock = new ReentrantReadWriteLock(); myFilePageCache = filePageCache; myUseReadWriteLock = useReadWriteLock; myCacheChannels = cacheChannels; + myDisableAssertions = disableAssertions; } public StorageLockContext(boolean useReadWriteLock, boolean cacheChannels) { - this(ourDefaultCache, useReadWriteLock, cacheChannels); + this(ourDefaultCache, useReadWriteLock, cacheChannels, false); + } + + public StorageLockContext(boolean useReadWriteLock, + boolean cacheChannels, + boolean disableAssertions) { + this(ourDefaultCache, useReadWriteLock, cacheChannels, disableAssertions); } public StorageLockContext(boolean useReadWriteLock) { - this(ourDefaultCache, useReadWriteLock, false); + this(ourDefaultCache, useReadWriteLock, false, false); } public StorageLockContext() { - this(ourDefaultCache, false, false); + this(ourDefaultCache, false, false, false); } boolean useChannelCache() { @@ -86,7 +95,7 @@ public final class StorageLockContext { @ApiStatus.Internal public void checkWriteAccess() { - if (IndexDebugProperties.DEBUG) { + if (!myDisableAssertions && IndexDebugProperties.DEBUG) { if (myLock.writeLock().isHeldByCurrentThread()) return; throw new IllegalStateException("Must hold StorageLock write lock to access PagedFileStorage"); } @@ -94,7 +103,7 @@ public final class StorageLockContext { @ApiStatus.Internal public void checkReadAccess() { - if (IndexDebugProperties.DEBUG) { + if (!myDisableAssertions && IndexDebugProperties.DEBUG) { if (myLock.getReadHoldCount() > 0 || myLock.writeLock().isHeldByCurrentThread()) return; throw new IllegalStateException("Must hold StorageLock read lock to access PagedFileStorage"); } } diff --git a/platform/vfs-impl/src/com/intellij/openapi/vfs/PlatformVirtualFileManager.java b/platform/vfs-impl/src/com/intellij/openapi/vfs/PlatformVirtualFileManager.java index 73e063c1a5c6..482e990cc575 100644 --- a/platform/vfs-impl/src/com/intellij/openapi/vfs/PlatformVirtualFileManager.java +++ b/platform/vfs-impl/src/com/intellij/openapi/vfs/PlatformVirtualFileManager.java @@ -40,7 +40,7 @@ public class PlatformVirtualFileManager extends VirtualFileManagerImpl { @Override public long getModificationCount() { - return myManagingFS.getModificationCount(); + return myManagingFS.getFilesystemModificationCount(); } @Override diff --git a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/FSRecords.java b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/FSRecords.java index e85205e67c1c..3ac69aded87e 100644 --- a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/FSRecords.java +++ b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/FSRecords.java @@ -79,7 +79,7 @@ public final class FSRecords { } private static int calculateVersion() { - return nextMask(59, // acceptable range is [0..255] + return nextMask(59 + (PersistentFSRecordsStorage.useLockFreeRecordsStorage ? 1 : 0), // acceptable range is [0..255] 8, nextMask(useContentHashes, nextMask(IOUtil.useNativeByteOrderForByteBuffers(), @@ -107,17 +107,19 @@ public final class FSRecords { /** * @return nameId > 0 */ - static int writeAttributesToRecord(int fileId, int parentId, @NotNull FileAttributes attributes, @NotNull String name, - boolean overwriteMissed) { + static int writeAttributesToRecord(int fileId, int parentId, @NotNull FileAttributes attributes, @NotNull String name, boolean overwriteMissed) { int nameId = getNameId(name); long timestamp = attributes.lastModified; long length = attributes.isDirectory() ? -1L : attributes.length; int flags = PersistentFSImpl.fileAttributesToFlags(attributes); - writeAndHandleErrors(() -> { + try { setAttributes(fileId, timestamp, length, flags, nameId, parentId, overwriteMissed); - return null; - }); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } InvertedNameIndex.updateFileName(fileId, nameId, 0); ourNamesIndexModCount.incrementAndGet(); @@ -130,9 +132,9 @@ public final class FSRecords { invalidateCaches(); int parentId = getParent(fileId); String msg = "File already created: fileId=" + fileId + - "; nameId=" + nameId + "(" + getNameByNameId(nameId) + ")" + + "; nameId=" + nameId + "(" + FileNameCache.getVFileName(nameId) + ")" + "; parentId=" + parentId + - "; existingData=" + existingData; + "; \nexistingData=" + existingData; if (parentId > 0) { msg += "; parent.name=" + getName(parentId); msg += "; parent.children=" + list(parentId); @@ -162,6 +164,7 @@ public final class FSRecords { ourTreeAccessor.ensureLoaded(); } catch (IOException e) { + LOG.error(e); handleError(e); } } @@ -177,7 +180,13 @@ public final class FSRecords { } public static long getCreationTimestamp() { - return readAndHandleErrors(() -> ourConnection.getTimestamp()); + try { + return ourConnection.getTimestamp(); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } // todo: Address / capacity store in records table, size store with payload @@ -218,27 +227,51 @@ public final class FSRecords { @TestOnly static int @NotNull [] listRoots() { - return readAndHandleErrors(() -> ourTreeAccessor.listRoots()); + try { + return ourTreeAccessor.listRoots(); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } @TestOnly static void force() { - writeAndHandleErrors(ourConnection::doForce); + try { + ourConnection.doForce(); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } @TestOnly static boolean isDirty() { - return readAndHandleErrors(ourConnection::isDirty); + return ourConnection.isDirty(); } @PersistentFS.Attributes static int getFlags(int id) { - return readAndHandleErrors(() -> ourConnection.getRecords().doGetFlags(id)); + try { + return ourConnection.getRecords().getFlags(id); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } @ApiStatus.Internal public static boolean isDeleted(int id) { - return readAndHandleErrors(() -> ourRecordAccessor.isDeleted(id)); + try { + return ourRecordAccessor.isDeleted(id); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } static int findRootRecord(@NotNull String rootUrl) { @@ -258,17 +291,35 @@ public final class FSRecords { } public static int @NotNull [] listIds(int fileId) { - return readAndHandleErrors(() -> ourTreeAccessor.listIds(fileId)); + try { + return ourTreeAccessor.listIds(fileId); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } static boolean mayHaveChildren(int fileId) { - return readAndHandleErrors(() -> ourTreeAccessor.mayHaveChildren(fileId)); + try { + return ourTreeAccessor.mayHaveChildren(fileId); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } // returns child infos (sorted by id) without (potentially expensive) name (or without even nameId if `loadNameId` is false) @NotNull static ListResult list(int parentId) { - return readAndHandleErrors(() -> ourTreeAccessor.doLoadChildren(parentId)); + try { + return ourTreeAccessor.doLoadChildren(parentId); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } @NotNull @@ -277,10 +328,16 @@ public final class FSRecords { } static boolean wereChildrenAccessed(int id) { - return readAndHandleErrors(() -> ourTreeAccessor.wereChildrenAccessed(id)); + try { + return ourTreeAccessor.wereChildrenAccessed(id); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } - static T readAndHandleErrors(@NotNull ThrowableComputable action) { + static T read(@NotNull ThrowableComputable action) { // otherwise DbConnection.handleError(e) (requires write lock) could fail if (lock.getReadHoldCount() != 0) { try { @@ -304,6 +361,19 @@ public final class FSRecords { // long reads like processXXX can be safely cancelled throw e; } + catch (Throwable e) { + throw new RuntimeException(e); + } + } + + static T readAndHandleErrors(@NotNull ThrowableComputable action) { + try { + return read(action); + } + catch (ProcessCanceledException e) { + // long reads like processXXX can be safely cancelled + throw e; + } catch (Throwable e) { handleError(e); throw new RuntimeException(e); @@ -454,13 +524,17 @@ public final class FSRecords { } static @Nullable String readSymlinkTarget(int id) { - String result = readAndHandleErrors(() -> { - try (DataInputStream stream = readAttribute(id, ourSymlinkTargetAttr)) { - if (stream != null) return StringUtil.nullize(IOUtil.readUTF(stream)); + try (DataInputStream stream = readAttribute(id, ourSymlinkTargetAttr)) { + if (stream != null) { + String result = StringUtil.nullize(IOUtil.readUTF(stream)); + return result == null ? null : FileUtil.toSystemIndependentName(result); } - return null; - }); - return result != null ? FileUtil.toSystemIndependentName(result) : null; + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } + return null; } static void storeSymlinkTarget(int id, @Nullable String symlinkTarget) { @@ -479,23 +553,31 @@ public final class FSRecords { @TestOnly static int getPersistentModCount() { - return readAndHandleErrors(ourConnection::getPersistentModCount); + return ourConnection.getPersistentModCount(); } - private static void incModCount(int id) throws IOException { + static void incModCount(int id) throws IOException { ourConnection.incModCount(id); } + public static int incGlobalModCount() throws IOException { + return ourConnection.incGlobalModCount(); + } + public static int getParent(int id) { - return readAndHandleErrors(() -> { - final int parentId = ourConnection.getRecords().getParent(id); + try { + int parentId = ourConnection.getRecords().getParent(id); if (parentId == id) { LOG.error("Cyclic parent child relations in the database. id = " + id); return 0; } return parentId; - }); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } @Nullable @@ -569,7 +651,7 @@ public final class FSRecords { @ApiStatus.Internal @NotNull public static IntList getRemainFreeRecords() { - return readAndHandleErrors(() -> new IntArrayList(ourConnection.getFreeRecords())); + return ourConnection.getFreeRecords(); } @ApiStatus.Internal @@ -584,14 +666,23 @@ public final class FSRecords { return; } - writeAndHandleErrors(() -> { - incModCount(id); + try { + incGlobalModCount(); ourConnection.getRecords().setParent(id, parentId); - }); + } catch (Throwable e) { + handleError(e); + throw new RuntimeException(e); + } } public static boolean processAllNames(@NotNull Processor processor) { - return readAndHandleErrors(() -> ourConnection.getNames().processAllDataObjects(processor)); + try { + return ourConnection.getNames().processAllDataObjects(processor); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } public static boolean processFilesWithNames(@NotNull Set names, @NotNull IntPredicate processor) { @@ -618,7 +709,14 @@ public final class FSRecords { @NotNull static CharSequence getNameSequence(int id) { - int nameId = readAndHandleErrors(() -> ourConnection.getRecords().getNameId(id)); + int nameId = 0; + try { + nameId = ourConnection.getRecords().getNameId(id); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } return nameId == 0 ? "" : FileNameCache.getVFileName(nameId); } @@ -652,45 +750,67 @@ public final class FSRecords { } static void setFlags(int id, @PersistentFS.Attributes int flags) { - writeAndHandleErrors(() -> { - incModCount(id); + try { ourConnection.getRecords().setFlags(id, flags); - }); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } static long getLength(int id) { - return readAndHandleErrors(() -> ourConnection.getRecords().getLength(id)); + try { + return ourConnection.getRecords().getLength(id); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } static void setLength(int id, long len) { - writeAndHandleErrors(() -> { - PersistentFSRecordsStorage records = ourConnection.getRecords(); - if (records.putLength(id, len)) { - incModCount(id); - } - }); + try { + ourConnection.getRecords().putLength(id, len); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } static void setAttributes(int id, long timestamp, long length, int flags, int nameId, int parentId, boolean overwriteMissed) throws IOException { - PersistentFSRecordsStorage records = ourConnection.getRecords(); - records.setAttributesAndIncModCount(id, timestamp, length, flags, nameId, parentId, overwriteMissed); + ourConnection.getRecords().setAttributesAndIncModCount(id, timestamp, length, flags, nameId, parentId, overwriteMissed); } static long getTimestamp(int id) { - return readAndHandleErrors(() -> ourConnection.getRecords().getTimestamp(id)); + try { + return ourConnection.getRecords().getTimestamp(id); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } static void setTimestamp(int id, long value) { - writeAndHandleErrors(() -> { - PersistentFSRecordsStorage records = ourConnection.getRecords(); - if (records.putTimeStamp(id, value)) { - incModCount(id); - } - }); + try { + ourConnection.getRecords().putTimestamp(id, value); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } static int getModCount(int id) { - return readAndHandleErrors(() -> ourConnection.getRecords().getModCount(id)); + try { + return ourConnection.getRecords().getModCount(id); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } @Nullable @@ -759,12 +879,24 @@ public final class FSRecords { } static int getContentId(int fileId) { - return readAndHandleErrors(() -> ourConnection.getRecords().getContentRecordId(fileId)); + try { + return ourConnection.getRecords().getContentRecordId(fileId); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } @TestOnly static byte[] getContentHash(int fileId) { - return readAndHandleErrors(() -> ourContentAccessor.getContentHash(fileId)); + try { + return ourContentAccessor.getContentHash(fileId); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } @NotNull diff --git a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSConnection.java b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSConnection.java index 5532f0fd2138..4d572aefd6c7 100644 --- a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSConnection.java +++ b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSConnection.java @@ -17,6 +17,7 @@ import com.intellij.util.io.storage.CapacityAllocationPolicy; import com.intellij.util.io.storage.HeavyProcessLatch; import com.intellij.util.io.storage.RefCountingContentStorage; import com.intellij.util.io.storage.Storage; +import it.unimi.dsi.fastutil.ints.IntArrayList; import it.unimi.dsi.fastutil.ints.IntList; import org.jetbrains.annotations.Contract; import org.jetbrains.annotations.NotNull; @@ -137,7 +138,9 @@ final class PersistentFSConnection { @NotNull IntList getFreeRecords() { - return myFreeRecords; + synchronized (myFreeRecords) { + return new IntArrayList(myFreeRecords); + } } long getTimestamp() throws IOException { @@ -175,7 +178,7 @@ final class PersistentFSConnection { return myRecords.getGlobalModCount(); } - int incGlobalModCount() throws IOException { + public int incGlobalModCount() throws IOException { incLocalModCount(); return myRecords.incGlobalModCount(); } @@ -189,7 +192,6 @@ final class PersistentFSConnection { void incLocalModCount() throws IOException { markDirty(); - //noinspection NonAtomicOperationOnVolatileField myLocalModificationCount.incrementAndGet(); } @@ -212,10 +214,13 @@ final class PersistentFSConnection { // must not be run under write lock to avoid other clients wait for read lock private void flush() { if (isDirty() && !HeavyProcessLatch.INSTANCE.isRunning()) { - FSRecords.readAndHandleErrors(() -> { + try { doForce(); - return null; - }); + } + catch (IOException e) { + handleError(e); + throw new RuntimeException(e); + } } } @@ -273,12 +278,10 @@ final class PersistentFSConnection { } } - // either called from FlushingDaemon thread under read lock, or from handleError under write lock void markClean() throws IOException { - assert FSRecords.lock.isWriteLocked() || FSRecords.lock.getReadHoldCount() != 0; + // no synchronization, it's ok to have race here if (myDirty) { myDirty = false; - // writing here under read lock is safe because no-one else read or write at this offset (except at startup) myRecords.setConnectionStatus(myCorrupted ? PersistentFSHeaders.CORRUPTED_MAGIC : PersistentFSHeaders.SAFELY_CLOSED_MAGIC); diff --git a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSConnector.java b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSConnector.java index 76fed45b2004..d14aacabd5f5 100644 --- a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSConnector.java +++ b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSConnector.java @@ -30,7 +30,7 @@ final class PersistentFSConnector { private static final int MAX_INITIALIZATION_ATTEMPTS = 10; private static final AtomicInteger INITIALIZATION_COUNTER = new AtomicInteger(); private static final StorageLockContext PERSISTENT_FS_STORAGE_CONTEXT = new StorageLockContext(false, true); - private static final StorageLockContext PERSISTENT_FS_STORAGE_CONTEXT_RW = new StorageLockContext(true, true); + private static final StorageLockContext PERSISTENT_FS_STORAGE_CONTEXT_RW = new StorageLockContext(true, true, true); static @NotNull PersistentFSConnection connect(@NotNull String cachesDir, int version, boolean useContentHashes) { return FSRecords.writeAndHandleErrors(() -> { @@ -117,16 +117,17 @@ final class PersistentFSConnector { SimpleStringPersistentEnumerator enumeratedAttributes = new SimpleStringPersistentEnumerator(enumeratedAttributesFile); - boolean aligned = PagedFileStorage.BUFFER_SIZE % PersistentFSRecordsStorage.RECORD_SIZE == 0; + int pageSize = PagedFileStorage.BUFFER_SIZE * PersistentFSRecordsStorage.recordsLength() / PersistentFSSynchronizedRecordsStorage.RECORD_SIZE; + boolean aligned = pageSize % PersistentFSRecordsStorage.recordsLength() == 0; if (!aligned) { - LOG.error("Buffer size " + PagedFileStorage.BUFFER_SIZE + " is not aligned for record size " + PersistentFSRecordsStorage.RECORD_SIZE); + LOG.error("Buffer size " + PagedFileStorage.BUFFER_SIZE + " is not aligned for record size " + PersistentFSRecordsStorage.recordsLength()); } - records = new PersistentFSRecordsStorage(new ResizeableMappedFile(recordsFile, - 20 * 1024, - PERSISTENT_FS_STORAGE_CONTEXT_RW, - PagedFileStorage.BUFFER_SIZE, - aligned, - IOUtil.useNativeByteOrderForByteBuffers())); + records = PersistentFSRecordsStorage.createStorage(new ResizeableMappedFile(recordsFile, + 20 * 1024, + PERSISTENT_FS_STORAGE_CONTEXT_RW, + pageSize, + aligned, + IOUtil.useNativeByteOrderForByteBuffers())); boolean initial = records.length() == 0; diff --git a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSHeaders.java b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSHeaders.java index dde6261c8884..5ad3eb7a0156 100644 --- a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSHeaders.java +++ b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSHeaders.java @@ -15,6 +15,7 @@ final class PersistentFSHeaders { static { //noinspection ConstantConditions - assert HEADER_SIZE <= PersistentFSRecordsStorage.RECORD_SIZE; + assert HEADER_SIZE <= PersistentFSLockFreeRecordsStorage.RECORD_SIZE; + assert HEADER_SIZE <= PersistentFSSynchronizedRecordsStorage.RECORD_SIZE; } } diff --git a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSImpl.java b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSImpl.java index efc7cb1f5ae2..625cc1be7d65 100644 --- a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSImpl.java +++ b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSImpl.java @@ -302,7 +302,7 @@ public final class PersistentFSImpl extends PersistentFS implements Disposable { @Override public int getModificationCount() { - return FSRecords.getLocalModCount(); + return FSRecords.getPersistentModCount(); } @Override diff --git a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSLockFreeRecordsStorage.kt b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSLockFreeRecordsStorage.kt new file mode 100644 index 000000000000..382421648986 --- /dev/null +++ b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSLockFreeRecordsStorage.kt @@ -0,0 +1,389 @@ +// Copyright 2000-2022 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. +package com.intellij.openapi.vfs.newvfs.persistent + +import com.intellij.util.containers.Unsafe +import com.intellij.util.io.ByteBufferUtil +import com.intellij.util.io.DirectByteBufferAllocator +import com.intellij.util.io.ResizeableMappedFile +import org.jetbrains.annotations.ApiStatus +import java.io.IOException +import java.nio.ByteBuffer +import java.nio.ByteOrder +import java.nio.channels.ReadableByteChannel +import java.util.concurrent.atomic.AtomicInteger +import kotlin.concurrent.withLock +import kotlin.math.max + +@ApiStatus.Internal +internal class PersistentFSLockFreeRecordsStorage @Throws(IOException::class) constructor(private val file: ResizeableMappedFile): PersistentFSRecordsStorage { + private val metadataReadLock = file.storageLockContext.readLock() + private val metadataWriteLock = file.storageLockContext.writeLock() + + private val _globalModCount: AtomicInteger + private val recordCount: AtomicInteger + + init { + _globalModCount = AtomicInteger(readGlobalModCount()) + recordCount = AtomicInteger((length() / RECORD_SIZE).toInt()) + } + + override fun getGlobalModCount(): Int = _globalModCount.get() + + @Throws(IOException::class) + private fun readGlobalModCount(): Int = metadataReadLock.withLock { + file.getInt(PersistentFSHeaders.HEADER_GLOBAL_MOD_COUNT_OFFSET.toLong()) + } + + @Throws(IOException::class) + private fun saveGlobalModCount() = metadataWriteLock.withLock { + file.putInt(PersistentFSHeaders.HEADER_GLOBAL_MOD_COUNT_OFFSET.toLong(), globalModCount) + } + + override fun incGlobalModCount(): Int { + return _globalModCount.incrementAndGet() + } + + @Throws(IOException::class) + override fun getTimestamp(): Long = metadataReadLock.withLock { + file.getLong(PersistentFSHeaders.HEADER_TIMESTAMP_OFFSET.toLong()) + } + + @Throws(IOException::class) + override fun setVersion(version: Int) = metadataWriteLock.withLock { + file.putInt(PersistentFSHeaders.HEADER_VERSION_OFFSET.toLong(), version) + file.putLong(PersistentFSHeaders.HEADER_TIMESTAMP_OFFSET.toLong(), System.currentTimeMillis()) + } + + @Throws(IOException::class) + override fun getVersion(): Int = metadataReadLock.withLock { + file.getInt(PersistentFSHeaders.HEADER_VERSION_OFFSET.toLong()) + } + + @Throws(IOException::class) + override fun getConnectionStatus() = metadataReadLock.withLock { + file.getInt(PersistentFSHeaders.HEADER_CONNECTION_STATUS_OFFSET.toLong()) + } + + @Throws(IOException::class) + override fun setConnectionStatus(code: Int) = metadataWriteLock.withLock { + file.putInt(PersistentFSHeaders.HEADER_CONNECTION_STATUS_OFFSET.toLong(), code) + } + + @Throws(IOException::class) + override fun getNameId(id: Int): Int = acquireRecord(id, AccessType.READ) { + nameId() + } + + @Throws(IOException::class) + override fun setNameId(id: Int, nameId: Int) = acquireRecord(id, AccessType.WRITE) { + nameId(nameId) + } + + @Throws(IOException::class) + override fun getParent(id: Int): Int = acquireRecord(id, AccessType.READ) { + parent() + } + + @Throws(IOException::class) + override fun setParent(id: Int, parent: Int) = acquireRecord(id, AccessType.WRITE) { + parent(parent) + } + + @Throws(IOException::class) + override fun getModCount(id: Int): Int = acquireRecord(id, AccessType.READ) { + modCount() + } + + @Throws(IOException::class) + override fun getFlags(id: Int): @PersistentFS.Attributes Int = acquireRecord(id, AccessType.READ) { + flags() + } + + @Throws(IOException::class) + override fun setFlags(id: Int, flags: @PersistentFS.Attributes Int) = acquireRecord(id, AccessType.WRITE_AND_INCREMENT_MOD_COUNTER) { + flags(flags) + } + + @Throws(IOException::class) + override fun setModCount(id: Int, value: Int) = acquireRecord(id, AccessType.WRITE) { + //TODO + incModCount() + } + + @Throws(IOException::class) + override fun getContentRecordId(id: Int): Int = acquireRecord(id, AccessType.READ) { + content() + } + + @Throws(IOException::class) + override fun setContentRecordId(id: Int, value: Int) = acquireRecord(id, AccessType.WRITE) { + content(value) + } + + @Throws(IOException::class) + override fun getAttributeRecordId(id: Int): Int = acquireRecord(id, AccessType.READ) { + attrs() + } + + @Throws(IOException::class) + override fun setAttributeRecordId(id: Int, value: Int) = acquireRecord(id, AccessType.WRITE) { + attrs(value) + } + + @Throws(IOException::class) + override fun getTimestamp(id: Int): Long = acquireRecord(id, AccessType.READ) { + timeStamp() + } + + @Throws(IOException::class) + override fun putTimestamp(id: Int, value: Long) = acquireRecord(id, AccessType.WRITE_AND_INCREMENT_MOD_COUNTER) { + timeStamp(value) + } + + @Throws(IOException::class) + override fun getLength(id: Int): Long = acquireRecord(id, AccessType.READ) { + length() + } + + @Throws(IOException::class) + override fun putLength(id: Int, value: Long) = acquireRecord(id, AccessType.WRITE_AND_INCREMENT_MOD_COUNTER) { + length(value) + } + + @Throws(IOException::class) + override fun cleanRecord(id: Int) = metadataWriteLock.withLock { + recordCount.updateAndGet { operand: Int -> max(id + 1, operand) } + file.put(id.toLong() * RECORD_SIZE, ZEROES, 0, RECORD_SIZE) + } + + override fun allocateRecord(): Int = recordCount.getAndIncrement() + + @Throws(IOException::class) + override fun setAttributesAndIncModCount(id: Int, + timestamp: Long, + length: Long, + flags: Int, + nameId: Int, + parentId: Int, + overwriteMissed: Boolean) = acquireRecord(id, AccessType.WRITE) { + setup(parentId, nameId, flags, 0, 0, timestamp, length, overwriteMissed) + } + + override fun length(): Long = metadataReadLock.withLock { + file.length() + } + + @Throws(IOException::class) + override fun close() = metadataWriteLock.withLock { + saveGlobalModCount() + file.close() + } + + @Throws(IOException::class) + override fun force() = metadataWriteLock.withLock { + saveGlobalModCount() + file.force() + } + + override fun isDirty(): Boolean = file.isDirty + + @Throws(IOException::class) + override fun processAllNames(operator: PersistentFSRecordsStorage.NameFlagsProcessor) = this.metadataReadLock.withLock { + // skip header + file.force() + // skip header + file.readChannel { ch: ReadableByteChannel -> + val buffer = ByteBuffer.allocateDirect(RECORD_SIZE * 1024) + try { + var id = 1 + var limit: Int + var offset: Int + while (ch.read(buffer).also { limit = it } >= RECORD_SIZE) { + offset = if (id == 1) RECORD_SIZE else 0 // skip header + while (offset < limit) { + val nameId = buffer.getInt(offset + NAME_OFFSET) + val flags = buffer.getInt(offset + FLAGS_OFFSET) + operator.process(id, nameId, flags) + id++ + offset += RECORD_SIZE + } + buffer.position(0) + } + } + catch (ignore: IOException) { + } + true + } + } + + private enum class AccessType { + READ, + WRITE, + WRITE_AND_INCREMENT_MOD_COUNTER + } + + private inline fun acquireRecord(id: Int, access: AccessType, action: LockFreeRecord.() -> V): V { + val recordOffset = id.toLong() * RECORD_SIZE + val storagePage = file.pagedFileStorage.getByteBuffer(recordOffset, false) + val inPageOffset = file.pagedFileStorage.getOffsetInPage(recordOffset) + + try { + if (access != AccessType.READ) { + storagePage.markDirty() + } + + val buffer = storagePage.buffer + val recordBuffer = buffer.duplicate().order(buffer.order()).limit(inPageOffset + RECORD_SIZE).position(inPageOffset).mark().slice() + return action(LockFreeRecord(recordBuffer)) + } + finally { + if (access != AccessType.READ) { + if (access == AccessType.WRITE_AND_INCREMENT_MOD_COUNTER) { + incGlobalModCount() + } + storagePage.fileSizeMayChanged(inPageOffset + RECORD_SIZE) + file.logicalSize = max(file.logicalSize, recordOffset + RECORD_SIZE) + } + storagePage.unlock() + } + } + + companion object { + private const val PARENT_OFFSET = 0 + private const val PARENT_SIZE = 4 + private const val NAME_OFFSET = PARENT_OFFSET + PARENT_SIZE + private const val NAME_SIZE = 4 + private const val FLAGS_OFFSET = NAME_OFFSET + NAME_SIZE + private const val FLAGS_SIZE = 4 + private const val ATTR_REF_OFFSET = FLAGS_OFFSET + FLAGS_SIZE + private const val ATTR_REF_SIZE = 4 + private const val CONTENT_OFFSET = ATTR_REF_OFFSET + ATTR_REF_SIZE + private const val CONTENT_SIZE = 4 + private const val TIMESTAMP_OFFSET = CONTENT_OFFSET + CONTENT_SIZE + private const val TIMESTAMP_SIZE = 8 + private const val LENGTH_OFFSET = TIMESTAMP_OFFSET + TIMESTAMP_SIZE + private const val LENGTH_SIZE = 8 + private const val MOD_COUNT_SIZE = 4 + private const val MOD_COUNT_PRE_OFFSET = LENGTH_OFFSET + LENGTH_SIZE + private const val MOD_COUNT_AFTER_OFFSET = MOD_COUNT_PRE_OFFSET + MOD_COUNT_SIZE + + const val RECORD_SIZE = MOD_COUNT_AFTER_OFFSET + MOD_COUNT_SIZE + private val ZEROES = ByteArray(RECORD_SIZE) + } + + @JvmInline + private value class LockFreeRecord(val data: ByteBuffer) { + init { + val bufferSize = data.limit() - data.position() + assert(bufferSize == RECORD_SIZE) { + "buffer size = $bufferSize" + } + } + + fun setup(parent: Int, nameId: Int, flags: Int, attrs: Int, content: Int, timestamp: Long, length: Long, overwriteMissed: Boolean) = update { + putInt(PARENT_OFFSET, parent) + putInt(NAME_OFFSET, nameId) + putInt(FLAGS_OFFSET, flags) + if (overwriteMissed) { + putInt(ATTR_REF_OFFSET, attrs) + } + putInt(CONTENT_OFFSET, content) + putLong(TIMESTAMP_OFFSET, timestamp) + putLong(LENGTH_OFFSET, length) + } + + fun parent(): Int = read { getInt(PARENT_OFFSET) } + fun parent(value: Int) = update { putInt(PARENT_OFFSET, value) } + + fun nameId(): Int = read { getInt(NAME_OFFSET) } + fun nameId(value: Int) = update { putInt(NAME_OFFSET, value) } + + fun flags(): Int = read { getInt(FLAGS_OFFSET) } + fun flags(value: Int) = update { putInt(FLAGS_OFFSET, value) } + + fun attrs(): Int = read { getInt(ATTR_REF_OFFSET) } + fun attrs(value: Int) = update { putInt(ATTR_REF_OFFSET, value) } + + fun content(): Int = read { getInt(CONTENT_OFFSET) } + fun content(value: Int) = update { putInt(CONTENT_OFFSET, value) } + + fun timeStamp(): Long = read { getLong(TIMESTAMP_OFFSET) } + fun timeStamp(value: Long) = update { putLong(TIMESTAMP_OFFSET, value) } + + fun length(): Long = read { getLong(LENGTH_OFFSET) } + fun length(value: Long) = update { putLong(LENGTH_OFFSET, value) } + + fun modCount() = modCountPre() + fun incModCount() = update { /*No Op*/ } + + fun modCountPre(): Int = data.getInt(MOD_COUNT_PRE_OFFSET) + fun modCountAfter(): Int = data.getInt(MOD_COUNT_AFTER_OFFSET) + + inline fun update(updater: ByteBuffer.() -> Unit) { + val writeBuffer = DirectByteBufferAllocator.ALLOCATOR.allocate(RECORD_SIZE) + writeBuffer.order(data.order()) + + try { + while (true) { + val readModCount = readTo(writeBuffer) + + updater(writeBuffer) + writeBuffer.rewind() + + if (tryWrite(writeBuffer, readModCount)) { + return + } + } + } + finally { + DirectByteBufferAllocator.ALLOCATOR.release(writeBuffer) + } + } + + fun readTo(copy: ByteBuffer): Int { + while (true) { + val modCountPre = modCountPre() + + copy.rewind() + copy.put(data.duplicate().order(data.order())) + + if (modCountAfter() == modCountPre) { + return modCountPre + } + } + } + + private inline fun read(eval: ByteBuffer.() -> V): V { + while (true) { + val modCountPre = modCountPre() + + val result = eval(data) + + val modCountAfter = modCountAfter() + if (modCountAfter == modCountPre) { + return result + } + } + } + + fun tryWrite(buffer: ByteBuffer, expectedModCounter: Int): Boolean { + val newModCounter = expectedModCounter + 1 + if (data.compareAndSwapInt(MOD_COUNT_PRE_OFFSET, expectedModCounter, newModCounter)) { + data.duplicate().order(data.order()).put(buffer.duplicate().order(data.order()).limit(RECORD_SIZE - 2 * MOD_COUNT_SIZE)) + data.putInt(MOD_COUNT_AFTER_OFFSET, newModCounter) + return true + } + return false + } + + companion object { + private fun ByteBuffer.compareAndSwapInt(offset: Int, expected: Int, new: Int): Boolean { + val address = ByteBufferUtil.getAddress(this) + val flipOrder = order() == ByteOrder.BIG_ENDIAN // TODO revise + val _expected = if (flipOrder) Integer.reverseBytes(expected) else expected + val _new = if (flipOrder) Integer.reverseBytes(new) else new + return Unsafe.compareAndSwapInt(null, address + offset, _expected, _new) + } + } + } +} \ No newline at end of file diff --git a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSRecordAccessor.java b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSRecordAccessor.java index a401f6b59c5b..53b70a2d6904 100644 --- a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSRecordAccessor.java +++ b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSRecordAccessor.java @@ -12,8 +12,6 @@ import org.jetbrains.annotations.Nullable; import java.io.IOException; -import static com.intellij.openapi.vfs.newvfs.persistent.PersistentFSRecordsStorage.RECORD_SIZE; - final class PersistentFSRecordAccessor { private static final Logger LOG = Logger.getInstance(PersistentFSRecordAccessor.class); static final int FREE_RECORD_FLAG = 0x400; @@ -72,13 +70,13 @@ final class PersistentFSRecordAccessor { long t = System.currentTimeMillis(); final int fileLength = length(); - assert fileLength % RECORD_SIZE == 0; - int recordCount = fileLength / RECORD_SIZE; + assert fileLength % PersistentFSRecordsStorage.recordsLength() == 0; + int recordCount = fileLength / PersistentFSRecordsStorage.recordsLength(); IntList usedAttributeRecordIds = new IntArrayList(); IntList validAttributeIds = new IntArrayList(); for (int id = 2; id < recordCount; id++) { - int flags = connection.getRecords().doGetFlags(id); + int flags = connection.getRecords().getFlags(id); LOG.assertTrue((flags & ~ALL_VALID_FLAGS) == 0, "Invalid flags: 0x" + Integer.toHexString(flags) + ", id: " + id); boolean isFreeRecord = connection.getFreeRecords().contains(id); if (BitUtil.isSet(flags, FREE_RECORD_FLAG)) { @@ -95,7 +93,7 @@ final class PersistentFSRecordAccessor { } boolean isDeleted(int id) throws IOException { - return BitUtil.isSet(myFSConnection.getRecords().doGetFlags(id), FREE_RECORD_FLAG) || myNewFreeRecords.contains(id); + return BitUtil.isSet(myFSConnection.getRecords().getFlags(id), FREE_RECORD_FLAG) || myNewFreeRecords.contains(id); } private void checkRecordSanity(int id, @@ -106,7 +104,7 @@ final class PersistentFSRecordAccessor { int parentId = connection.getRecords().getParent(id); assert parentId >= 0 && parentId < recordCount; if (parentId > 0 && connection.getRecords().getParent(parentId) > 0) { - int parentFlags = connection.getRecords().doGetFlags(parentId); + int parentFlags = connection.getRecords().getFlags(parentId); assert !BitUtil.isSet(parentFlags, FREE_RECORD_FLAG) : parentId + ": " + Integer.toHexString(parentFlags); assert BitUtil.isSet(parentFlags, PersistentFS.Flags.IS_DIRECTORY) : parentId + ": " + Integer.toHexString(parentFlags); } @@ -136,13 +134,10 @@ final class PersistentFSRecordAccessor { return (int)myFSConnection.getRecords().length(); } - private int allocateRecord() { return myFSConnection.getRecords().allocateRecord(); } - - private void deleteContentAndAttributes(int id) throws IOException { myPersistentFSContentAccessor.deleteContent(id); myPersistentFSAttributeAccessor.deleteAttributes(id); diff --git a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSRecordsStorage.java b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSRecordsStorage.java index 3b571f1c83f2..86adc372410d 100644 --- a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSRecordsStorage.java +++ b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSRecordsStorage.java @@ -1,338 +1,86 @@ -// Copyright 2000-2021 JetBrains s.r.o. Use of this source code is governed by the Apache 2.0 license that can be found in the LICENSE file. +// Copyright 2000-2022 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. package com.intellij.openapi.vfs.newvfs.persistent; -import com.intellij.openapi.util.ThrowableComputable; +import com.intellij.util.SystemProperties; import com.intellij.util.io.ResizeableMappedFile; +import org.jetbrains.annotations.ApiStatus; import org.jetbrains.annotations.NotNull; import java.io.IOException; -import java.nio.ByteBuffer; -import java.nio.ByteOrder; -import java.util.concurrent.atomic.AtomicInteger; -public final class PersistentFSRecordsStorage { - private static final int PARENT_OFFSET = 0; - private static final int PARENT_SIZE = 4; - private static final int NAME_OFFSET = PARENT_OFFSET + PARENT_SIZE; - private static final int NAME_SIZE = 4; - private static final int FLAGS_OFFSET = NAME_OFFSET + NAME_SIZE; - private static final int FLAGS_SIZE = 4; - private static final int ATTR_REF_OFFSET = FLAGS_OFFSET + FLAGS_SIZE; - private static final int ATTR_REF_SIZE = 4; - private static final int CONTENT_OFFSET = ATTR_REF_OFFSET + ATTR_REF_SIZE; - private static final int CONTENT_SIZE = 4; - private static final int TIMESTAMP_OFFSET = CONTENT_OFFSET + CONTENT_SIZE; - private static final int TIMESTAMP_SIZE = 8; - private static final int MOD_COUNT_OFFSET = TIMESTAMP_OFFSET + TIMESTAMP_SIZE; - private static final int MOD_COUNT_SIZE = 4; - private static final int LENGTH_OFFSET = MOD_COUNT_OFFSET + MOD_COUNT_SIZE; - private static final int LENGTH_SIZE = 8; +@ApiStatus.Internal +interface PersistentFSRecordsStorage { + boolean useLockFreeRecordsStorage = SystemProperties.getBooleanProperty("idea.use.lock.free.record.storage.for.vfs", false); - static final int RECORD_SIZE = LENGTH_OFFSET + LENGTH_SIZE; - private static final byte[] ZEROES = new byte[RECORD_SIZE]; - - private V read(ThrowableComputable action) throws E { - myFile.getStorageLockContext().lockRead(); - try { - return action.compute(); - } - finally { - myFile.getStorageLockContext().unlockRead(); - } + static int recordsLength() { + return useLockFreeRecordsStorage ? PersistentFSLockFreeRecordsStorage.RECORD_SIZE : PersistentFSSynchronizedRecordsStorage.RECORD_SIZE; } - private V write(ThrowableComputable action) throws E { - myFile.getStorageLockContext().lockWrite(); - try { - return action.compute(); - } - finally { - myFile.getStorageLockContext().unlockWrite(); - } + static PersistentFSRecordsStorage createStorage(@NotNull ResizeableMappedFile file) throws IOException { + return useLockFreeRecordsStorage ? new PersistentFSLockFreeRecordsStorage(file) : new PersistentFSSynchronizedRecordsStorage(file); } - @NotNull - private final ResizeableMappedFile myFile; - private final ByteBuffer myPooledWriteBuffer = ByteBuffer.allocateDirect(RECORD_SIZE); - @NotNull - private final AtomicInteger myGlobalModCount; - @NotNull - private final AtomicInteger myRecordCount; + int allocateRecord(); - public PersistentFSRecordsStorage(@NotNull ResizeableMappedFile file) throws IOException { - myFile = file; - if (myFile.isNativeBytesOrder()) myPooledWriteBuffer.order(ByteOrder.nativeOrder()); - myGlobalModCount = new AtomicInteger(readGlobalModCount()); - myRecordCount = new AtomicInteger((int)(length() / RECORD_SIZE)); - } + void setAttributeRecordId(int fileId, int recordId) throws IOException; - int getGlobalModCount() { - return myGlobalModCount.get(); - } + int getAttributeRecordId(int fileId) throws IOException; - private int readGlobalModCount() throws IOException { - return read(() -> { - return myFile.getInt(PersistentFSHeaders.HEADER_GLOBAL_MOD_COUNT_OFFSET); - }); - } + int getParent(int fileId) throws IOException; - private void saveGlobalModCount() throws IOException { - write(() -> { - myFile.putInt(PersistentFSHeaders.HEADER_GLOBAL_MOD_COUNT_OFFSET, getGlobalModCount()); - return null; - }); - } + void setParent(int fileIf, int parentId) throws IOException; - int incGlobalModCount() { - return myGlobalModCount.incrementAndGet(); - } + int getNameId(int fileId) throws IOException; - long getTimestamp() throws IOException { - return read(() -> { - return myFile.getLong(PersistentFSHeaders.HEADER_TIMESTAMP_OFFSET); - }); - } + void setNameId(int fileId, int nameId) throws IOException; - void setVersion(int version) throws IOException { - write(() -> { - myFile.putInt(PersistentFSHeaders.HEADER_VERSION_OFFSET, version); - myFile.putLong(PersistentFSHeaders.HEADER_TIMESTAMP_OFFSET, System.currentTimeMillis()); - return null; - }); - } + void setFlags(int fileId, int flags) throws IOException; - int getVersion() throws IOException { - return read(() -> { - return myFile.getInt(PersistentFSHeaders.HEADER_VERSION_OFFSET); - }); - } + long getLength(int fileId) throws IOException; - void setConnectionStatus(int connectionStatus) throws IOException { - write(() -> { - myFile.putInt(PersistentFSHeaders.HEADER_CONNECTION_STATUS_OFFSET, connectionStatus); - return null; - }); - } + void putLength(int fileId, long length) throws IOException; - int getConnectionStatus() throws IOException { - return read(() -> { - return myFile.getInt(PersistentFSHeaders.HEADER_CONNECTION_STATUS_OFFSET); - }); - } + long getTimestamp(int fileId) throws IOException; - int getNameId(int id) throws IOException { - return read(() -> { - assert id > 0 : id; - return getRecordInt(id, NAME_OFFSET); - }); - } + void putTimestamp(int fileId, long timestamp) throws IOException; - void setNameId(int id, int nameId) throws IOException { - write(() -> { - PersistentFSConnection.ensureIdIsValid(nameId); - putRecordInt(id, NAME_OFFSET, nameId); - return null; - }); - } + int getModCount(int fileId) throws IOException; - int getParent(int id) throws IOException { - return read(() -> { - return getRecordInt(id, PARENT_OFFSET); - }); - } + void setModCount(int fileId, int counter) throws IOException; - void setParent(int id, int parent) throws IOException { - write(() -> { - putRecordInt(id, PARENT_OFFSET, parent); - return null; - }); - } + int getContentRecordId(int fileId) throws IOException; - int getModCount(int id) throws IOException { - return read(() -> { - return getRecordInt(id, MOD_COUNT_OFFSET); - }); - } + void setContentRecordId(int fileId, int recordId) throws IOException; - @PersistentFS.Attributes - int doGetFlags(int id) throws IOException { - return read(() -> { - return getRecordInt(id, FLAGS_OFFSET); - }); - } + int getFlags(int fileId) throws IOException; - void setFlags(int id, @PersistentFS.Attributes int flags) throws IOException { - write(() -> { - putRecordInt(id, FLAGS_OFFSET, flags); - return null; - }); - } + void setAttributesAndIncModCount(int fileId, long timestamp, long length, int flags, int nameId, int parentId, boolean overwriteMissed) throws IOException; - void setModCount(int id, int value) throws IOException { - write(() -> { - putRecordInt(id, MOD_COUNT_OFFSET, value); - return null; - }); - } + boolean isDirty(); - int getContentRecordId(int fileId) throws IOException { - return read(() -> { - return getRecordInt(fileId, CONTENT_OFFSET); - }); - } + long getTimestamp() throws IOException; - void setContentRecordId(int id, int value) throws IOException { - write(() -> { - putRecordInt(id, CONTENT_OFFSET, value); - return null; - }); - } + void setConnectionStatus(int code) throws IOException; - int getAttributeRecordId(int id) throws IOException { - return read(() -> { - return getRecordInt(id, ATTR_REF_OFFSET); - }); - } + int getConnectionStatus() throws IOException; - void setAttributeRecordId(int id, int value) throws IOException { - write(() -> { - putRecordInt(id, ATTR_REF_OFFSET, value); - return null; - }); - } + void setVersion(int version) throws IOException; - long getTimestamp(int id) throws IOException { - return read(() -> { - return myFile.getLong(getOffset(id, TIMESTAMP_OFFSET)); - }); - } + int getVersion() throws IOException; - boolean putTimeStamp(int id, long value) throws IOException { - return write(() -> { - int timeStampOffset = getOffset(id, TIMESTAMP_OFFSET); - if (myFile.getLong(timeStampOffset) != value) { - myFile.putLong(timeStampOffset, value); - return true; - } - return false; - }); - } + int getGlobalModCount(); - long getLength(int id) throws IOException { - return read(() -> { - return myFile.getLong(getOffset(id, LENGTH_OFFSET)); - }); - } + int incGlobalModCount(); - boolean putLength(int id, long value) throws IOException { - return write(() -> { - int lengthOffset = getOffset(id, LENGTH_OFFSET); - if (myFile.getLong(lengthOffset) != value) { - myFile.putLong(lengthOffset, value); - return true; - } - return false; - }); - } + long length(); - void cleanRecord(int id) throws IOException { - write(() -> { - myRecordCount.updateAndGet(operand -> Math.max(id + 1, operand)); - myFile.put(((long)id) * RECORD_SIZE, ZEROES, 0, RECORD_SIZE); - return null; - }); - } + void cleanRecord(int fileId) throws IOException; - int allocateRecord() { - return myRecordCount.getAndIncrement(); - } + boolean processAllNames(@NotNull NameFlagsProcessor processor) throws IOException; - private int getRecordInt(int id, int offset) throws IOException { - return read(() -> { - return myFile.getInt(getOffset(id, offset)); - }); - } + void force() throws IOException; - private void putRecordInt(int id, int offset, int value) throws IOException { - write(() -> { - myFile.putInt(getOffset(id, offset), value); - return null; - }); - } - - public void setAttributesAndIncModCount(int id, long timestamp, long length, int flags, int nameId, int parentId, - boolean overwriteMissed) throws IOException { - write(() -> { - assert myPooledWriteBuffer.position() == 0; - myPooledWriteBuffer.putLong(TIMESTAMP_OFFSET, timestamp); - myPooledWriteBuffer.putInt(ATTR_REF_OFFSET, overwriteMissed ? 0 : getAttributeRecordId(id)); - myPooledWriteBuffer.putLong(LENGTH_OFFSET, length); - myPooledWriteBuffer.putInt(FLAGS_OFFSET, flags); - myPooledWriteBuffer.putInt(NAME_OFFSET, nameId); - myPooledWriteBuffer.putInt(PARENT_OFFSET, parentId); - myPooledWriteBuffer.putInt(MOD_COUNT_OFFSET, overwriteMissed ? 0 : getModCount(id) + 1); - assert myPooledWriteBuffer.position() == 0; - myFile.put(((long)id) * RECORD_SIZE, myPooledWriteBuffer); - myPooledWriteBuffer.rewind(); - return null; - }); - } - - private static int getOffset(int id, int offset) { - return id * RECORD_SIZE + offset; - } - - long length() { - return read(() -> { - return myFile.length(); - }); - } - - void close() throws IOException { - write(() -> { - saveGlobalModCount(); - myFile.close(); - return null; - }); - } - - void force() throws IOException { - write(() -> { - saveGlobalModCount(); - myFile.force(); - return null; - }); - } - - boolean isDirty() { - return myFile.isDirty(); - } - - void processAllNames(@NotNull NameFlagsProcessor operator) throws IOException { - read(() -> { - myFile.force(); - return myFile.readChannel(ch -> { - ByteBuffer buffer = ByteBuffer.allocateDirect(RECORD_SIZE * 1024); - if (myFile.isNativeBytesOrder()) buffer.order(ByteOrder.nativeOrder()); - try { - int id = 1, limit, offset; - while ((limit = ch.read(buffer)) >= RECORD_SIZE) { - offset = id == 1 ? RECORD_SIZE : 0; // skip header - for (; offset < limit; offset += RECORD_SIZE) { - int nameId = buffer.getInt(offset + NAME_OFFSET); - int flags = buffer.getInt(offset + FLAGS_OFFSET); - operator.process(id, nameId, flags); - id ++; - } - buffer.position(0); - } - } - catch (IOException ignore) { - } - return true; - }); - }); - } + void close() throws IOException; interface NameFlagsProcessor { void process(int fileId, int nameId, int flags); diff --git a/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSSynchronizedRecordsStorage.java b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSSynchronizedRecordsStorage.java new file mode 100644 index 000000000000..8f26a25569d8 --- /dev/null +++ b/platform/vfs-impl/src/com/intellij/openapi/vfs/newvfs/persistent/PersistentFSSynchronizedRecordsStorage.java @@ -0,0 +1,368 @@ +// Copyright 2000-2022 JetBrains s.r.o. and contributors. Use of this source code is governed by the Apache 2.0 license. +package com.intellij.openapi.vfs.newvfs.persistent; + +import com.intellij.openapi.util.ThrowableComputable; +import com.intellij.util.io.ResizeableMappedFile; +import org.jetbrains.annotations.ApiStatus; +import org.jetbrains.annotations.NotNull; + +import java.io.IOException; +import java.nio.ByteBuffer; +import java.nio.ByteOrder; +import java.util.concurrent.atomic.AtomicInteger; + +@ApiStatus.Internal +final class PersistentFSSynchronizedRecordsStorage implements PersistentFSRecordsStorage { + private static final int PARENT_OFFSET = 0; + private static final int PARENT_SIZE = 4; + private static final int NAME_OFFSET = PARENT_OFFSET + PARENT_SIZE; + private static final int NAME_SIZE = 4; + private static final int FLAGS_OFFSET = NAME_OFFSET + NAME_SIZE; + private static final int FLAGS_SIZE = 4; + private static final int ATTR_REF_OFFSET = FLAGS_OFFSET + FLAGS_SIZE; + private static final int ATTR_REF_SIZE = 4; + private static final int CONTENT_OFFSET = ATTR_REF_OFFSET + ATTR_REF_SIZE; + private static final int CONTENT_SIZE = 4; + private static final int TIMESTAMP_OFFSET = CONTENT_OFFSET + CONTENT_SIZE; + private static final int TIMESTAMP_SIZE = 8; + private static final int MOD_COUNT_OFFSET = TIMESTAMP_OFFSET + TIMESTAMP_SIZE; + private static final int MOD_COUNT_SIZE = 4; + private static final int LENGTH_OFFSET = MOD_COUNT_OFFSET + MOD_COUNT_SIZE; + private static final int LENGTH_SIZE = 8; + + static final int RECORD_SIZE = LENGTH_OFFSET + LENGTH_SIZE; + private static final byte[] ZEROES = new byte[RECORD_SIZE]; + + private V read(ThrowableComputable action) throws E { + myFile.getStorageLockContext().lockRead(); + try { + return action.compute(); + } + finally { + myFile.getStorageLockContext().unlockRead(); + } + } + + private V write(ThrowableComputable action) throws E { + myFile.getStorageLockContext().lockWrite(); + try { + return action.compute(); + } + finally { + myFile.getStorageLockContext().unlockWrite(); + } + } + + @NotNull + private final ResizeableMappedFile myFile; + private final ByteBuffer myPooledWriteBuffer = ByteBuffer.allocateDirect(RECORD_SIZE); + @NotNull + private final AtomicInteger myGlobalModCount; + @NotNull + private final AtomicInteger myRecordCount; + + PersistentFSSynchronizedRecordsStorage(@NotNull ResizeableMappedFile file) throws IOException { + myFile = file; + if (myFile.isNativeBytesOrder()) myPooledWriteBuffer.order(ByteOrder.nativeOrder()); + myGlobalModCount = new AtomicInteger(readGlobalModCount()); + myRecordCount = new AtomicInteger((int)(length() / RECORD_SIZE)); + } + + @Override + public int getGlobalModCount() { + return myGlobalModCount.get(); + } + + private int readGlobalModCount() throws IOException { + return read(() -> { + return myFile.getInt(PersistentFSHeaders.HEADER_GLOBAL_MOD_COUNT_OFFSET); + }); + } + + private void saveGlobalModCount() throws IOException { + write(() -> { + myFile.putInt(PersistentFSHeaders.HEADER_GLOBAL_MOD_COUNT_OFFSET, getGlobalModCount()); + return null; + }); + } + + @Override + public int incGlobalModCount() { + return myGlobalModCount.incrementAndGet(); + } + + @Override + public long getTimestamp() throws IOException { + return read(() -> { + return myFile.getLong(PersistentFSHeaders.HEADER_TIMESTAMP_OFFSET); + }); + } + + @Override + public void setVersion(int version) throws IOException { + write(() -> { + myFile.putInt(PersistentFSHeaders.HEADER_VERSION_OFFSET, version); + myFile.putLong(PersistentFSHeaders.HEADER_TIMESTAMP_OFFSET, System.currentTimeMillis()); + return null; + }); + } + + @Override + public int getVersion() throws IOException { + return read(() -> { + return myFile.getInt(PersistentFSHeaders.HEADER_VERSION_OFFSET); + }); + } + + @Override + public void setConnectionStatus(int connectionStatus) throws IOException { + write(() -> { + myFile.putInt(PersistentFSHeaders.HEADER_CONNECTION_STATUS_OFFSET, connectionStatus); + return null; + }); + } + + @Override + public int getConnectionStatus() throws IOException { + return read(() -> { + return myFile.getInt(PersistentFSHeaders.HEADER_CONNECTION_STATUS_OFFSET); + }); + } + + @Override + public int getNameId(int id) throws IOException { + return read(() -> { + assert id > 0 : id; + return getRecordInt(id, NAME_OFFSET); + }); + } + + @Override + public void setNameId(int id, int nameId) throws IOException { + write(() -> { + PersistentFSConnection.ensureIdIsValid(nameId); + putRecordInt(id, NAME_OFFSET, nameId); + return null; + }); + } + + @Override + public int getParent(int id) throws IOException { + return read(() -> { + return getRecordInt(id, PARENT_OFFSET); + }); + } + + @Override + public void setParent(int id, int parent) throws IOException { + write(() -> { + putRecordInt(id, PARENT_OFFSET, parent); + return null; + }); + } + + @Override + public int getModCount(int id) throws IOException { + return read(() -> { + return getRecordInt(id, MOD_COUNT_OFFSET); + }); + } + + @Override + @PersistentFS.Attributes + public int getFlags(int id) throws IOException { + return read(() -> { + return getRecordInt(id, FLAGS_OFFSET); + }); + } + + @Override + public void setFlags(int id, @PersistentFS.Attributes int flags) throws IOException { + write(() -> { + FSRecords.incModCount(id); + putRecordInt(id, FLAGS_OFFSET, flags); + return null; + }); + } + + @Override + public void setModCount(int id, int value) throws IOException { + write(() -> { + putRecordInt(id, MOD_COUNT_OFFSET, value); + return null; + }); + } + + @Override + public int getContentRecordId(int fileId) throws IOException { + return read(() -> { + return getRecordInt(fileId, CONTENT_OFFSET); + }); + } + + @Override + public void setContentRecordId(int id, int value) throws IOException { + write(() -> { + putRecordInt(id, CONTENT_OFFSET, value); + return null; + }); + } + + @Override + public int getAttributeRecordId(int id) throws IOException { + return read(() -> { + return getRecordInt(id, ATTR_REF_OFFSET); + }); + } + + @Override + public void setAttributeRecordId(int id, int value) throws IOException { + write(() -> { + putRecordInt(id, ATTR_REF_OFFSET, value); + return null; + }); + } + + @Override + public long getTimestamp(int id) throws IOException { + return read(() -> { + return myFile.getLong(getOffset(id, TIMESTAMP_OFFSET)); + }); + } + + @Override + public void putTimestamp(int id, long value) throws IOException { + write(() -> { + int timeStampOffset = getOffset(id, TIMESTAMP_OFFSET); + if (myFile.getLong(timeStampOffset) != value) { + myFile.putLong(timeStampOffset, value); + FSRecords.incModCount(id); + } + return null; + }); + } + + @Override + public long getLength(int id) throws IOException { + return read(() -> { + return myFile.getLong(getOffset(id, LENGTH_OFFSET)); + }); + } + + @Override + public void putLength(int id, long value) throws IOException { + write(() -> { + int lengthOffset = getOffset(id, LENGTH_OFFSET); + if (myFile.getLong(lengthOffset) != value) { + myFile.putLong(lengthOffset, value); + FSRecords.incModCount(id); + } + return null; + }); + } + + @Override + public void cleanRecord(int id) throws IOException { + write(() -> { + myRecordCount.updateAndGet(operand -> Math.max(id + 1, operand)); + myFile.put(((long)id) * RECORD_SIZE, ZEROES, 0, RECORD_SIZE); + return null; + }); + } + + @Override + public int allocateRecord() { + return myRecordCount.getAndIncrement(); + } + + private int getRecordInt(int id, int offset) throws IOException { + return read(() -> { + return myFile.getInt(getOffset(id, offset)); + }); + } + + private void putRecordInt(int id, int offset, int value) throws IOException { + write(() -> { + myFile.putInt(getOffset(id, offset), value); + return null; + }); + } + + @Override + public void setAttributesAndIncModCount(int id, long timestamp, long length, int flags, int nameId, int parentId, boolean overwriteMissed) throws IOException { + write(() -> { + assert myPooledWriteBuffer.position() == 0; + myPooledWriteBuffer.putLong(TIMESTAMP_OFFSET, timestamp); + myPooledWriteBuffer.putInt(ATTR_REF_OFFSET, overwriteMissed ? 0 : getAttributeRecordId(id)); + myPooledWriteBuffer.putLong(LENGTH_OFFSET, length); + myPooledWriteBuffer.putInt(FLAGS_OFFSET, flags); + myPooledWriteBuffer.putInt(NAME_OFFSET, nameId); + myPooledWriteBuffer.putInt(PARENT_OFFSET, parentId); + assert myPooledWriteBuffer.position() == 0; + myFile.put(((long)id) * RECORD_SIZE, myPooledWriteBuffer); + myPooledWriteBuffer.rewind(); + return null; + }); + } + + private static int getOffset(int id, int offset) { + return id * RECORD_SIZE + offset; + } + + @Override + public long length() { + return read(() -> { + return myFile.length(); + }); + } + + @Override + public void close() throws IOException { + write(() -> { + saveGlobalModCount(); + myFile.close(); + return null; + }); + } + + @Override + public void force() throws IOException { + write(() -> { + saveGlobalModCount(); + myFile.force(); + return null; + }); + } + + @Override + public boolean isDirty() { + return myFile.isDirty(); + } + + @Override + public boolean processAllNames(@NotNull NameFlagsProcessor operator) throws IOException { + return read(() -> { + myFile.force(); + return myFile.readChannel(ch -> { + ByteBuffer buffer = ByteBuffer.allocateDirect(RECORD_SIZE * 1024); + if (myFile.isNativeBytesOrder()) buffer.order(ByteOrder.nativeOrder()); + try { + int id = 1, limit, offset; + while ((limit = ch.read(buffer)) >= RECORD_SIZE) { + offset = id == 1 ? RECORD_SIZE : 0; // skip header + for (; offset < limit; offset += RECORD_SIZE) { + int nameId = buffer.getInt(offset + NAME_OFFSET); + int flags = buffer.getInt(offset + FLAGS_OFFSET); + operator.process(id, nameId, flags); + id ++; + } + buffer.position(0); + } + } + catch (IOException ignore) { + } + return true; + }); + }); + } +}