From 85abe7b3aa576f6d01686d014faab3098ccc7375 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Tue, 4 May 2021 15:29:07 +0300 Subject: [PATCH] Fix race in invalidating entries --- .../mledger/impl/EntryCacheImpl.java | 2 +- .../bookkeeper/mledger/impl/EntryImpl.java | 15 +++++++- .../bookkeeper/mledger/impl/OpAddEntry.java | 8 ++--- .../util/InvalidateableReferenceCounted.java | 32 +++++++++++++++++ .../bookkeeper/mledger/util/RangeCache.java | 36 +++++++++++-------- .../mledger/util/RangeCacheTest.java | 15 ++++++-- 6 files changed, 86 insertions(+), 22 deletions(-) create mode 100644 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/InvalidateableReferenceCounted.java diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheImpl.java index ee3660b9e2a86..90a412b0aaab5 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryCacheImpl.java @@ -116,7 +116,7 @@ public boolean insert(EntryImpl entry) { return true; } else { // entry was not inserted into cache, we need to discard it - cacheEntry.release(); + cacheEntry.invalidate(); return false; } } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java index e25b52ac8cd70..08a232247653f 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java @@ -27,11 +27,14 @@ import io.netty.util.Recycler.Handle; import io.netty.util.ReferenceCounted; +import java.util.concurrent.atomic.AtomicBoolean; import org.apache.bookkeeper.client.api.LedgerEntry; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.util.AbstractCASReferenceCounted; +import org.apache.bookkeeper.mledger.util.InvalidateableReferenceCounted; -public final class EntryImpl extends AbstractCASReferenceCounted implements Entry, Comparable, ReferenceCounted { +public final class EntryImpl extends AbstractCASReferenceCounted implements Entry, Comparable, + InvalidateableReferenceCounted { private static final Recycler RECYCLER = new Recycler() { @Override @@ -45,6 +48,7 @@ protected EntryImpl newObject(Handle handle) { private long ledgerId; private long entryId; ByteBuf data; + private final AtomicBoolean invalidated = new AtomicBoolean(); public static EntryImpl create(LedgerEntry ledgerEntry) { EntryImpl entry = RECYCLER.get(); @@ -166,7 +170,16 @@ protected void deallocate() { timestamp = -1; ledgerId = -1; entryId = -1; + invalidated.set(false); recyclerHandle.recycle(this); } + @Override + public boolean invalidate() { + if (invalidated.compareAndSet(false, true)) { + release(); + return true; + } + return false; + } } diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java index e30c578fb556d..350602cb99cf0 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/OpAddEntry.java @@ -154,7 +154,7 @@ public void addComplete(int rc, final LedgerHandle lh, long entryId, Object ctx) } checkArgument(ledger.getId() == lh.getId(), "ledgerId %s doesn't match with acked ledgerId %s", ledger.getId(), lh.getId()); - + if (!checkAndCompleteOp(ctx)) { // means callback might have been completed by different thread (timeout task thread).. so do nothing return; @@ -189,7 +189,7 @@ public void safeRun() { // EntryCache.insert: duplicates entry by allocating new entry and data. so, recycle entry after calling // insert ml.entryCache.insert(entry); - entry.release(); + entry.invalidate(); } PositionImpl lastEntry = PositionImpl.get(ledger.getId(), entryId); @@ -248,7 +248,7 @@ private void updateLatency() { /** * Checks if add-operation is completed - * + * * @return true if task is not already completed else returns false. */ private boolean checkAndCompleteOp(Object ctx) { @@ -269,7 +269,7 @@ void handleAddTimeoutFailure(final LedgerHandle ledger, Object ctx) { /** * It handles add failure on the given ledger. it can be triggered when add-entry fails or times out. - * + * * @param ledger */ void handleAddFailure(final LedgerHandle ledger) { diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/InvalidateableReferenceCounted.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/InvalidateableReferenceCounted.java new file mode 100644 index 0000000000000..f5b69086ca69a --- /dev/null +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/InvalidateableReferenceCounted.java @@ -0,0 +1,32 @@ +/** + * 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.bookkeeper.mledger.util; + +import io.netty.util.ReferenceCounted; + +public interface InvalidateableReferenceCounted extends ReferenceCounted { + /** + * Marks the entry to be removed. Internally {@link ReferenceCounted#release()} will be called unless + * the entry has already been invalidated before. No calls to {@link ReferenceCounted#release()} should be + * made separately to invalidate the entry. + * + * @return true if the value was marked to be invalidated, false if it was already invalidated before. + */ + boolean invalidate(); +} diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/RangeCache.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/RangeCache.java index a5786ad867034..6f8a8169983c1 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/RangeCache.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/util/RangeCache.java @@ -21,7 +21,7 @@ import static com.google.common.base.Preconditions.checkArgument; import com.google.common.collect.Lists; -import io.netty.util.ReferenceCounted; +import io.netty.util.IllegalReferenceCountException; import java.util.Collection; import java.util.List; import java.util.Map; @@ -38,7 +38,7 @@ * @param * Cache value */ -public class RangeCache, Value extends ReferenceCounted> { +public class RangeCache, Value extends InvalidateableReferenceCounted> { // Map from key to nodes inside the linked list private final ConcurrentNavigableMap entries; private AtomicLong size; // Total size of values stored in cache @@ -90,7 +90,7 @@ public Value get(Key key) { try { value.retain(); return value; - } catch (Throwable t) { + } catch (IllegalReferenceCountException e) { // Value was already destroyed between get() and retain() return null; } @@ -113,7 +113,7 @@ public Collection getRange(Key first, Key last) { try { value.retain(); values.add(value); - } catch (Throwable t) { + } catch (IllegalReferenceCountException e) { // Value was already destroyed between get() and retain() } } @@ -140,9 +140,11 @@ public Pair removeRange(Key first, Key last, boolean lastInclusiv continue; } - removedSize += weighter.getSize(value); - value.release(); - ++removedEntries; + long entrySize = weighter.getSize(value); + if (value.invalidate()) { + removedSize += entrySize; + ++removedEntries; + } } size.addAndGet(-removedSize); @@ -168,9 +170,11 @@ public Pair evictLeastAccessedEntries(long minSize) { } Value value = entry.getValue(); - ++removedEntries; - removedSize += weighter.getSize(value); - value.release(); + long entrySize = weighter.getSize(value); + if (value.invalidate()) { + ++removedEntries; + removedSize += entrySize; + } } size.addAndGet(-removedSize); @@ -197,8 +201,10 @@ public long evictLEntriesBeforeTimestamp(long maxTimestamp) { } Value value = entry.getValue(); - removedSize += weighter.getSize(value); - value.release(); + long entrySize = weighter.getSize(value); + if (value.invalidate()) { + removedSize += entrySize; + } } size.addAndGet(-removedSize); @@ -230,8 +236,10 @@ public synchronized long clear() { break; } Value value = entry.getValue(); - removedSize += weighter.getSize(value); - value.release(); + long entrySize = weighter.getSize(value); + if (value.invalidate()) { + removedSize += entrySize; + } } entries.clear(); diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/RangeCacheTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/RangeCacheTest.java index 95896d24f35f2..36dca8ddeb73e 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/RangeCacheTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/util/RangeCacheTest.java @@ -27,13 +27,15 @@ import com.google.common.collect.Lists; import io.netty.util.AbstractReferenceCounted; import io.netty.util.ReferenceCounted; +import java.util.concurrent.atomic.AtomicBoolean; import org.apache.commons.lang3.tuple.Pair; import org.testng.annotations.Test; public class RangeCacheTest { - class RefString extends AbstractReferenceCounted implements ReferenceCounted { + class RefString extends AbstractReferenceCounted implements InvalidateableReferenceCounted { final String s; + private final AtomicBoolean invalidated = new AtomicBoolean(); RefString(String s) { super(); @@ -43,7 +45,7 @@ class RefString extends AbstractReferenceCounted implements ReferenceCounted { @Override protected void deallocate() { - // no-op + invalidated.set(false); } @Override @@ -61,6 +63,15 @@ public boolean equals(Object obj) { return false; } + + @Override + public boolean invalidate() { + if (invalidated.compareAndSet(false, true)) { + release(); + return true; + } + return false; + } } @Test