From ec0689888b61193722314aa2dd5731897bdecbd1 Mon Sep 17 00:00:00 2001 From: "Maxim.Mossienko" Date: Sat, 9 Jun 2012 17:11:46 +0400 Subject: [PATCH] avoid deadlock + readContent can returned data of unfinished write --- .../util/io/storage/RefCountingStorage.java | 92 +++++++++++-------- 1 file changed, 52 insertions(+), 40 deletions(-) diff --git a/platform/util/src/com/intellij/util/io/storage/RefCountingStorage.java b/platform/util/src/com/intellij/util/io/storage/RefCountingStorage.java index 675ce786e2d6..b8704c5d2432 100644 --- a/platform/util/src/com/intellij/util/io/storage/RefCountingStorage.java +++ b/platform/util/src/com/intellij/util/io/storage/RefCountingStorage.java @@ -51,7 +51,7 @@ public class RefCountingStorage extends AbstractStorage { private final boolean myDoNotZipCaches = Boolean.valueOf(System.getProperty("idea.doNotZipCaches")).booleanValue(); private static final int MAX_PENDING_ZIP_SIZE = 20 * 1024 * 1024; - private final Map myPendingWriteRequests = new ConcurrentHashMap(); + private final Map myPendingWriteRequests = new ConcurrentHashMap(); private volatile int myPendingWriteRequestsSize; private final LowMemoryWatcher myPendingWritesFlusher = LowMemoryWatcher.register(new Runnable() { @Override @@ -60,19 +60,19 @@ public class RefCountingStorage extends AbstractStorage { } }); - //private static class WriteRequest { - // final byte[] content; - // final int length; - // final int recordId; - // final boolean fixedSize; - // - // WriteRequest(byte[] _content, int _length, int _recordId, boolean _fixedSize) { - // content = _content; - // length = _length; - // recordId = _recordId; - // fixedSize = _fixedSize; - // } - //} + private static class WriteRequest { + final byte[] content; + final int length; + final int recordId; + final boolean fixedSize; + + WriteRequest(byte[] _content, int _length, int _recordId, boolean _fixedSize) { + content = _content; + length = _length; + recordId = _recordId; + fixedSize = _fixedSize; + } + } private static final int MAX_PENDING_WRITE_SIZE = 2 * 1024 * 1024; @@ -93,10 +93,23 @@ public class RefCountingStorage extends AbstractStorage { } private BufferExposingByteArrayOutputStream internalReadStream(int record) throws IOException { - waitForPendingWriteForRecord(record); + waitForZipToFinish(record); + WriteRequest request; + synchronized (myLock) { + request = myPendingWriteRequests.get(record); + } - byte[] result = super.readBytes(record); - InflaterInputStream in = new CustomInflaterInputStream(result); + byte[] bytes; + int length; + if (request != null) { + bytes = request.content; + length = request.length; + } else { + bytes = super.readBytes(record); + length = bytes.length; + } + + InflaterInputStream in = new CustomInflaterInputStream(bytes, length); try { final BufferExposingByteArrayOutputStream outputStream = new BufferExposingByteArrayOutputStream(); StreamUtil.copyStreamContent(in, outputStream); @@ -108,17 +121,20 @@ public class RefCountingStorage extends AbstractStorage { } private static class CustomInflaterInputStream extends InflaterInputStream { - public CustomInflaterInputStream(byte[] compressedData) { - super(new UnsyncByteArrayInputStream(compressedData), new Inflater(), 1); + private int usedBufferLength; + + public CustomInflaterInputStream(byte[] compressedData, int _length) { + super(new UnsyncByteArrayInputStream(compressedData, 0, _length), new Inflater(), 1); // force to directly use compressed data, this ensures less round trips with native extraction code and copy streams this.buf = compressedData; this.len = -1; // ensure one time fill + usedBufferLength = _length; } @Override protected void fill() throws IOException { if (len >= 0) throw new EOFException(); - len = buf.length; + len = usedBufferLength; inf.setInput(buf, 0, len); } @@ -132,13 +148,13 @@ public class RefCountingStorage extends AbstractStorage { private void waitForPendingWriteForRecord(int record) { waitForZipToFinish(record); - Callable action; + WriteRequest request; synchronized (myLock) { - action = myPendingWriteRequests.get(record); + request = myPendingWriteRequests.get(record); } - if (action != null) { + if (request != null) { try { - action.call(); + write(request); } catch (Exception e) { throw new RuntimeException(e); @@ -176,8 +192,9 @@ public class RefCountingStorage extends AbstractStorage { return; } + waitForPendingWriteForRecord(record); // ensure previous write was completed + synchronized (myLock) { - waitForPendingWriteForRecord(record); // ensure previous write was completed myPendingZipRequestsSize += bytes.getLength(); if (myPendingZipRequestsSize > MAX_PENDING_ZIP_SIZE) { // help async thread @@ -201,13 +218,8 @@ public class RefCountingStorage extends AbstractStorage { private void scheduleZippedContentToWrite(final BufferExposingByteArrayOutputStream outputStream, final int record, final boolean fixedSize) { synchronized (myLock) { myPendingWriteRequestsSize += outputStream.size(); - myPendingWriteRequests.put(record, new Callable() { - @Override - public Void call() throws Exception { - write(outputStream, record, fixedSize); - return null; - } - }); + myPendingWriteRequests.put(record, new WriteRequest(outputStream.getInternalBuffer(), outputStream.size(), record, fixedSize)); + if (myPendingWriteRequestsSize > MAX_PENDING_WRITE_SIZE) { // we do it under lock to ensure normally only one thread will bulky flush stuff flushPendingWrites(); } @@ -215,12 +227,12 @@ public class RefCountingStorage extends AbstractStorage { } - private void write(BufferExposingByteArrayOutputStream zippedBytes, int record, boolean fixedSize) throws IOException { + private void write(WriteRequest writeRequest) throws IOException { synchronized (myLock) { - if (!myPendingWriteRequests.containsKey(record)) return; // some thread helped us - super.writeBytes(record, new ByteSequence(zippedBytes.getInternalBuffer(), 0, zippedBytes.size()), fixedSize); - myPendingWriteRequests.remove(record); - myPendingWriteRequestsSize -= zippedBytes.size(); + if (!myPendingWriteRequests.containsKey(writeRequest.recordId)) return; // some thread helped us + super.writeBytes(writeRequest.recordId, new ByteSequence(writeRequest.content, 0, writeRequest.length), writeRequest.fixedSize); + myPendingWriteRequests.remove(writeRequest.recordId); + myPendingWriteRequestsSize -= writeRequest.length; } } @@ -306,10 +318,10 @@ public class RefCountingStorage extends AbstractStorage { } private void flushPendingWrites() { - for(Map.Entry entry: myPendingWriteRequests.entrySet()) { + for(Map.Entry entry: myPendingWriteRequests.entrySet()) { try { - Callable value = entry.getValue(); - if (value != null) value.call(); + WriteRequest value = entry.getValue(); + if (value != null) write(value); } catch (Exception e) { throw new RuntimeException(e); }