Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -172,33 +172,34 @@ public Manifest getManifest() throws IOException {
@Override
public Enumeration<JarEntry> entries() {
synchronized (this) {
ensureOpen();
return new JarEntriesEnumeration(this.resources.zipContent());
ZipContent zipContent = ensureOpen();
return new JarEntriesEnumeration(zipContent);
}
}

@Override
public Stream<JarEntry> stream() {
synchronized (this) {
ensureOpen();
return streamContentEntries().map(NestedJarEntry::new);
ZipContent zipContent = ensureOpen();
return streamContentEntries(zipContent).map(NestedJarEntry::new);
}
}

@Override
public Stream<JarEntry> versionedStream() {
synchronized (this) {
ensureOpen();
return streamContentEntries().map(this::getBaseName)
ZipContent zipContent = ensureOpen();
return streamContentEntries(zipContent).map(this::getBaseName)
.filter(Objects::nonNull)
.distinct()
.map(this::getJarEntry)
.filter(Objects::nonNull);
}

}

private Stream<ZipContent.Entry> streamContentEntries() {
ZipContentEntriesSpliterator spliterator = new ZipContentEntriesSpliterator(this.resources.zipContent());
private Stream<ZipContent.Entry> streamContentEntries(ZipContent zipContent) {
ZipContentEntriesSpliterator spliterator = new ZipContentEntriesSpliterator(zipContent);
return StreamSupport.stream(spliterator, false);
}

Expand Down Expand Up @@ -248,10 +249,8 @@ public boolean hasEntry(String name) {
if (entry != null) {
return true;
}
synchronized (this) {
ensureOpen();
return this.resources.zipContent().hasEntry(null, name);
}
ZipContent zipContent = ensureOpen();
return zipContent.hasEntry(null, name);
}

private NestedJarEntry getNestedJarEntry(String name) {
Expand All @@ -260,13 +259,17 @@ private NestedJarEntry getNestedJarEntry(String name) {
if (lastEntry != null && name.equals(lastEntry.getName())) {
return lastEntry;
}
ZipContent.Entry entry = getVersionedContentEntry(name);
entry = (entry != null) ? entry : getContentEntry(null, name);
if (entry == null) {
return null;
NestedJarEntry nestedJarEntry;
synchronized (this) {
ZipContent.Entry entry = getVersionedContentEntry(name);
entry = (entry != null) ? entry : getContentEntry(null, name);

if (entry == null) {
return null;
}
nestedJarEntry = new NestedJarEntry(entry, name);
this.lastEntry = nestedJarEntry;
}
NestedJarEntry nestedJarEntry = new NestedJarEntry(entry, name);
this.lastEntry = nestedJarEntry;
return nestedJarEntry;
}

Expand All @@ -291,10 +294,8 @@ private ZipContent.Entry getVersionedContentEntry(String name) {
}

private ZipContent.Entry getContentEntry(String namePrefix, String name) {
synchronized (this) {
ensureOpen();
return this.resources.zipContent().getEntry(namePrefix, name);
}
ZipContent zipContent = ensureOpen();
return zipContent.getEntry(namePrefix, name);
}

private ManifestInfo getManifestInfo() {
Expand All @@ -303,8 +304,8 @@ private ManifestInfo getManifestInfo() {
return manifestInfo;
}
synchronized (this) {
ensureOpen();
manifestInfo = this.resources.zipContent().getInfo(ManifestInfo.class, this::getManifestInfo);
ZipContent zipContent = ensureOpen();
manifestInfo = zipContent.getInfo(ManifestInfo.class, this::getManifestInfo);
}
this.manifestInfo = manifestInfo;
return manifestInfo;
Expand All @@ -331,11 +332,8 @@ private MetaInfVersionsInfo getMetaInfVersionsInfo() {
if (metaInfVersionsInfo != null) {
return metaInfVersionsInfo;
}
synchronized (this) {
ensureOpen();
metaInfVersionsInfo = this.resources.zipContent()
.getInfo(MetaInfVersionsInfo.class, MetaInfVersionsInfo::get);
}
ZipContent zipContent = ensureOpen();
metaInfVersionsInfo = zipContent.getInfo(MetaInfVersionsInfo.class, MetaInfVersionsInfo::get);
this.metaInfVersionsInfo = metaInfVersionsInfo;
return metaInfVersionsInfo;
}
Expand Down Expand Up @@ -374,16 +372,16 @@ private InputStream getInputStream(ZipContent.Entry contentEntry) throws IOExcep
@Override
public String getComment() {
synchronized (this) {
ensureOpen();
return this.resources.zipContent().getComment();
ZipContent zipContent = ensureOpen();
return zipContent.getComment();
}
}

@Override
public int size() {
synchronized (this) {
ensureOpen();
return this.resources.zipContent().size();
synchronized (this) { // consistent with superclass.
ZipContent zipContent = ensureOpen();
return zipContent.size();
}
}

Expand All @@ -409,13 +407,21 @@ public String getName() {
return this.name;
}

private void ensureOpen() {
/**
* Ensures the jar is open and that there is a {@link ZipContent} instance available.
* @return the ZipContent instance - never {@code null}.
* @throws IllegalStateException if the jar is closed or the {@link ZipContent} is
* null.
*/
private ZipContent ensureOpen() {
if (this.closed) {
throw new IllegalStateException("Zip file closed");
}
if (this.resources.zipContent() == null) {
ZipContent zipContent = this.resources.zipContent();

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

While validating this approach I realised that the zipContent variable being returned from the zipContent() method isn't a volatile field so there is a small gap here. Adding volatile to this field in NestedJarFileResources would close the gap. Without changing it to volatile though the code should still be safe as the reference counting in the FileDataBlock would catch it and throw a consistent error (no corruption or deadlock).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added additional concurrency tests to verify the safety of this change as it stands. These tests are what found the bug in NestedJarFileResources. The test (NestedJarFileConcurrencyTest) was intended to confirm that the reference count check in FileDataBlock is sufficient to make a stale read of the non-volatile zipContent field fail cleanly.

I have been unable to reproduce the scenario of a stale read, but the possible difference is that a ClosedChannelException is thrown where an IllegalStateException was previously thrown. I did consider catching it and rethrowing as IllegalStateException, but consider the scenario unlikely enough that I have not.

if (zipContent == null) {
throw new IllegalStateException("The object is not initialized.");
}
return zipContent;
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@
* support for slicing.
*
* @author Phillip Webb
* @author Ian Kettle

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not sure if i should add this - changes to the concurrency feel significant enough to meet threshold in the contribution doc. If its added here it should add to the NestedJarFile change too.

*/
class FileDataBlock implements CloseableDataBlock {

Expand Down Expand Up @@ -81,7 +82,7 @@ public int read(ByteBuffer dst, long pos) throws IOException {
long updatedLimit = dst.position() + remaining;
dst.limit((updatedLimit > Integer.MAX_VALUE) ? Integer.MAX_VALUE : (int) updatedLimit);
}
int result = this.fileAccess.read(dst, this.offset + pos);
int result = this.fileAccess.read(dst, this.offset + pos, ClosedChannelException::new);
if (originalDestinationLimit != -1) {
dst.limit(originalDestinationLimit);
}
Expand Down Expand Up @@ -160,7 +161,7 @@ static class FileAccess {

private final Path path;

private int referenceCount;
private volatile int referenceCount;

private FileChannel fileChannel;

Expand All @@ -183,8 +184,12 @@ static class FileAccess {
this.path = path;
}

int read(ByteBuffer dst, long position) throws IOException {
int read(ByteBuffer dst, long position, Supplier<? extends IOException> closedExceptionSupplier)
throws IOException {
synchronized (this.lock) {
if (this.referenceCount == 0) {
throw closedExceptionSupplier.get();
}
if (position < this.bufferPosition || position >= this.bufferPosition + this.bufferSize) {
fillBuffer(position);
}
Expand Down Expand Up @@ -243,24 +248,26 @@ private void repairFileChannel() throws IOException {

void open() throws IOException {
synchronized (this.lock) {
if (this.referenceCount == 0) {
int localReferenceCount = this.referenceCount;
if (localReferenceCount == 0) {
debug.log("Opening '%s'", this.path);
this.fileChannel = FileChannel.open(this.path, StandardOpenOption.READ);
this.buffer = ByteBuffer.allocateDirect(BUFFER_SIZE);
tracker.openedFileChannel(this.path);
}
this.referenceCount++;
debug.log("Reference count for '%s' incremented to %s", this.path, this.referenceCount);
this.referenceCount = ++localReferenceCount;
debug.log("Reference count for '%s' incremented to %s", this.path, localReferenceCount);
}
}

void close() throws IOException {
synchronized (this.lock) {
if (this.referenceCount == 0) {
int localReferenceCount = this.referenceCount;
if (localReferenceCount == 0) {
return;
}
this.referenceCount--;
if (this.referenceCount == 0) {
this.referenceCount = --localReferenceCount;
if (localReferenceCount == 0) {
debug.log("Closing '%s'", this.path);
this.buffer = null;
this.bufferPosition = -1;
Expand All @@ -274,15 +281,13 @@ void close() throws IOException {
this.randomAccessFile = null;
}
}
debug.log("Reference count for '%s' decremented to %s", this.path, this.referenceCount);
debug.log("Reference count for '%s' decremented to %s", this.path, localReferenceCount);
}
}

<E extends Exception> void ensureOpen(Supplier<E> exceptionSupplier) throws E {
synchronized (this.lock) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I will need to find or reproduce the thread dump to confirm but I believe this was where the blocking moved to after sorting the NestedJarFile concurrency. There were 4 places inside the class competing for the same lock. The open, close, read and ensureOpen. The only usage of ensureOpen is in the same method call as the read. This change removes the synchronisation entirely from ensureOpen by using an AtomicInteger for the reference tracking which reduces internal contention. Previously FileDataBlock#read was needing to synchronize on the lock twice for each call.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Changed approach - reverted my original change as i didn't want to include it but then found that there was a gap in the read method meaning that the ensureOpen result could be stale by the time the read actually occurred and the read didn't recheck that it was still open in the sync block. Have moved referenceCount to be volatil so the ensureOpen no longer needs a sync block as it is informative only and the read now also checks the referenceCount within the same sync block that actually does the read

if (this.referenceCount == 0) {
throw exceptionSupplier.get();
}
if (this.referenceCount == 0) {
throw exceptionSupplier.get();
}
}

Expand Down
Loading
Loading