• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

ben-manes / caffeine / #5784

15 Sep 2026 02:57AM UTC coverage: 99.848% (-0.02%) from 99.869%
#5784

push

github

ben-manes
keep jcache access extensions across expiry races

A JCache read extends an entry's access expiry in two lock-free writes,
to the native timer and to the wrapper, and "Reschedule jcache access
expiry outside the compute" had put the native timer first. A read that
judged its captured wrapper expired removed it by identity alone, so it
could discard an entry another read had just extended and publish
EXPIRED for it, and any core read between the two writes re-derived the
native deadline from the older wrapper, letting the entry expire early.
The wrapper is written first again, and the expired-read removals in
containsKey, get, getAll and the loading get now require the captured
wrapper to still be expired.

4742 of 4802 branches covered (98.75%)

13 of 13 new or added lines in 2 files covered. (100.0%)

6 existing lines in 4 files now uncovered.

9211 of 9225 relevant lines covered (99.85%)

1.0 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

99.96
/caffeine/src/main/java/com/github/benmanes/caffeine/cache/BoundedLocalCache.java
1
/*
2
 * Copyright 2014 Ben Manes. All Rights Reserved.
3
 *
4
 * Licensed under the Apache License, Version 2.0 (the "License");
5
 * you may not use this file except in compliance with the License.
6
 * You may obtain a copy of the License at
7
 *
8
 *     http://www.apache.org/licenses/LICENSE-2.0
9
 *
10
 * Unless required by applicable law or agreed to in writing, software
11
 * distributed under the License is distributed on an "AS IS" BASIS,
12
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13
 * See the License for the specific language governing permissions and
14
 * limitations under the License.
15
 */
16
package com.github.benmanes.caffeine.cache;
17

18
import static com.github.benmanes.caffeine.cache.Async.ASYNC_EXPIRY;
19
import static com.github.benmanes.caffeine.cache.Caffeine.calculateHashMapCapacity;
20
import static com.github.benmanes.caffeine.cache.Caffeine.ceilingPowerOfTwo;
21
import static com.github.benmanes.caffeine.cache.Caffeine.requireArgument;
22
import static com.github.benmanes.caffeine.cache.Caffeine.requireState;
23
import static com.github.benmanes.caffeine.cache.Caffeine.toNanosSaturated;
24
import static com.github.benmanes.caffeine.cache.Caffeine.toUnchecked;
25
import static com.github.benmanes.caffeine.cache.LocalCache.castNonNull;
26
import static com.github.benmanes.caffeine.cache.LocalCache.nullRef;
27
import static com.github.benmanes.caffeine.cache.LocalLoadingCache.newBulkMappingFunction;
28
import static com.github.benmanes.caffeine.cache.LocalLoadingCache.newMappingFunction;
29
import static com.github.benmanes.caffeine.cache.Node.PROBATION;
30
import static com.github.benmanes.caffeine.cache.Node.PROTECTED;
31
import static com.github.benmanes.caffeine.cache.Node.WINDOW;
32
import static java.lang.invoke.ConstantBootstraps.fieldVarHandle;
33
import static java.util.Locale.US;
34
import static java.util.Objects.requireNonNull;
35
import static java.util.Spliterator.DISTINCT;
36
import static java.util.Spliterator.IMMUTABLE;
37
import static java.util.Spliterator.NONNULL;
38
import static java.util.Spliterator.ORDERED;
39

40
import java.io.InvalidObjectException;
41
import java.io.ObjectInputStream;
42
import java.io.Serializable;
43
import java.lang.System.Logger;
44
import java.lang.System.Logger.Level;
45
import java.lang.invoke.MethodHandles;
46
import java.lang.invoke.VarHandle;
47
import java.lang.ref.ReferenceQueue;
48
import java.lang.ref.WeakReference;
49
import java.time.Duration;
50
import java.util.AbstractCollection;
51
import java.util.AbstractSet;
52
import java.util.ArrayDeque;
53
import java.util.Collection;
54
import java.util.Collections;
55
import java.util.Comparator;
56
import java.util.Deque;
57
import java.util.HashMap;
58
import java.util.IdentityHashMap;
59
import java.util.Iterator;
60
import java.util.LinkedHashMap;
61
import java.util.Map;
62
import java.util.NoSuchElementException;
63
import java.util.Objects;
64
import java.util.Optional;
65
import java.util.OptionalInt;
66
import java.util.OptionalLong;
67
import java.util.Set;
68
import java.util.Spliterator;
69
import java.util.Spliterators;
70
import java.util.concurrent.CancellationException;
71
import java.util.concurrent.CompletableFuture;
72
import java.util.concurrent.ConcurrentHashMap;
73
import java.util.concurrent.ConcurrentMap;
74
import java.util.concurrent.Executor;
75
import java.util.concurrent.ForkJoinPool;
76
import java.util.concurrent.ForkJoinTask;
77
import java.util.concurrent.ThreadLocalRandom;
78
import java.util.concurrent.TimeUnit;
79
import java.util.concurrent.TimeoutException;
80
import java.util.concurrent.locks.ReentrantLock;
81
import java.util.function.BiConsumer;
82
import java.util.function.BiFunction;
83
import java.util.function.Consumer;
84
import java.util.function.Function;
85
import java.util.function.Predicate;
86
import java.util.stream.Stream;
87
import java.util.stream.StreamSupport;
88

89
import org.jspecify.annotations.NonNull;
90
import org.jspecify.annotations.Nullable;
91

92
import com.github.benmanes.caffeine.cache.Async.AsyncExpiry;
93
import com.github.benmanes.caffeine.cache.LinkedDeque.PeekingIterator;
94
import com.github.benmanes.caffeine.cache.Policy.CacheEntry;
95
import com.github.benmanes.caffeine.cache.References.InternalReference;
96
import com.github.benmanes.caffeine.cache.stats.StatsCounter;
97
import com.google.errorprone.annotations.CanIgnoreReturnValue;
98
import com.google.errorprone.annotations.Var;
99
import com.google.errorprone.annotations.concurrent.GuardedBy;
100

101
/**
102
 * An in-memory cache implementation that supports full concurrency of retrievals, a high expected
103
 * concurrency for updates, and multiple ways to bound the cache.
104
 * <p>
105
 * This class is abstract and code generated subclasses provide the complete implementation for a
106
 * particular configuration. This is to ensure that only the fields and execution paths necessary
107
 * for a given configuration are used.
108
 *
109
 * @author ben.manes@gmail.com (Ben Manes)
110
 * @param <K> the type of keys maintained by this cache
111
 * @param <V> the type of mapped values
112
 */
113
@SuppressWarnings({"RedundantSuppression", "ResultOfMethodCallIgnored", "serial", "unused"})
114
abstract class BoundedLocalCache<K, V> extends BLCHeader.DrainStatusRef
115
    implements LocalCache<K, V> {
116

117
  /*
118
   * This class performs a best-effort bounding of a ConcurrentHashMap using a page-replacement
119
   * algorithm to determine which entries to evict when the capacity is exceeded.
120
   *
121
   * Concurrency:
122
   * ------------
123
   * The page replacement algorithms are kept eventually consistent with the map. An update to the
124
   * map and recording of reads may not be immediately reflected in the policy's data structures.
125
   * These structures are guarded by a lock, and operations are applied in batches to avoid lock
126
   * contention. The penalty of applying the batches is spread across threads, so that the amortized
127
   * cost is slightly higher than performing just the ConcurrentHashMap operation [1].
128
   *
129
   * A memento of the reads and writes that were performed on the map is recorded in buffers. These
130
   * buffers are drained at the first opportunity after a write or when a read buffer is full. The
131
   * reads are offered to a buffer that will reject additions if contended on or if it is full. Due
132
   * to the concurrent nature of the read and write operations, a strict policy ordering is not
133
   * possible, but it may be observably strict when single-threaded. The buffers are drained
134
   * asynchronously to minimize the request latency and uses a state machine to determine when to
135
   * schedule this work on an executor.
136
   *
137
   * Due to a lack of a strict ordering guarantee, a task can be executed out-of-order, such as a
138
   * removal followed by its addition. The state of the entry is encoded using the key field to
139
   * avoid additional memory usage. An entry is "alive" if it is in both the hash table and the page
140
   * replacement policy. It is "retired" if it is not in the hash table and is pending removal from
141
   * the page replacement policy. Finally, an entry transitions to the "dead" state when it is
142
   * neither in the hash table nor the page replacement policy. Both the retired and dead states are
143
   * represented by a sentinel key that should not be used for map operations.
144
   *
145
   * Eviction:
146
   * ---------
147
   * Maximum size is implemented using the Window TinyLfu policy [2] due to its high hit rate, O(1)
148
   * time complexity, and small footprint. A new entry starts in the admission window and remains
149
   * there as long as it has high temporal locality (recency). Eventually an entry will slip from
150
   * the window into the main space. If the main space is already full, then a historic frequency
151
   * filter determines whether to evict the newly admitted entry or the victim entry chosen by the
152
   * eviction policy. This process ensures that the entries in the window were very recently used,
153
   * while entries in the main space are accessed very frequently and remain moderately recent. The
154
   * windowing allows the policy to have a high hit rate when entries exhibit a bursty access
155
   * pattern, while the filter ensures that popular items are retained. The admission window uses
156
   * LRU and the main space uses Segmented LRU.
157
   *
158
   * The optimal size of the window vs. main spaces is workload dependent [3]. A large admission
159
   * window is favored by recency-biased workloads, while a small one favors frequency-biased
160
   * workloads. When the window is too small, then recent arrivals are prematurely evicted, but when
161
   * it is too large, then they pollute the cache and force the eviction of more popular entries.
162
   * The optimal configuration is dynamically determined by using hill climbing to walk the hit
163
   * rate curve. For modestly sized caches this samples the hit rate and adjusts the window size
164
   * in the direction that is improving, decaying the step until the climber converges and
165
   * restarting when a large change indicates that the workload altered. Tiny caches allow the
166
   * sample period to grow as its step decays to avoid getting stuck at a poor configuration.
167
   * Larger caches instead compare the regions' hit densities within a single sample, which is
168
   * immune to the hit rate swings of a phasey workload, and take a proportional step towards the
169
   * balance point. As that signal only observes the hits that the current split earns, a region
170
   * that earns almost nothing cannot measure what a different split would, so those samples are
171
   * never trusted to hold position: the climber probes out of the blind state by hit rate,
172
   * accepts the new position if the density signal validates it, and otherwise undoes the probe
173
   * with exponentially longer pauses between attempts.
174
   *
175
   * The historic usage is retained in a compact popularity sketch, which uses hashing to
176
   * probabilistically estimate an item's frequency. This exposes a flaw where an adversary could
177
   * use hash flooding [4] to artificially raise the frequency of the main space's victim and cause
178
   * all candidates to be rejected. In the worst case, by exploiting hash collisions, an attacker
179
   * could cause the cache to never hit and hold only worthless items, resulting in a
180
   * denial-of-service attack against the underlying resource. This is mitigated by introducing
181
   * jitter, allowing candidates that are at least moderately popular to have a small, random chance
182
   * of being admitted. This causes the victim to be evicted, but in a way that marginally impacts
183
   * the hit rate.
184
   *
185
   * Expiration:
186
   * -----------
187
   * Expiration is implemented in O(1) time complexity. The time-to-idle policy uses an access-order
188
   * queue, the time-to-live policy uses a write-order queue, and variable expiration uses a
189
   * hierarchical timer wheel [5]. The queuing policies allow for peeking at the oldest entry to
190
   * determine if it has expired. If it has not, then the younger entries must not have expired
191
   * either. If a maximum size is set, then expiration will share the queues, minimizing the
192
   * per-entry footprint. The timer wheel based policy uses hashing and cascading in a manner that
193
   * amortizes the penalty of sorting to achieve a similar algorithmic cost.
194
   *
195
   * The expiration updates are applied in a best effort fashion. The reordering of variable or
196
   * access-order expiration may be discarded by the read buffer if it is full or contended.
197
   * Similarly, recording the touch for expiration to extend its lifetime may be ignored for an
198
   * entry if the last update was within a short time window. This is done to avoid overwhelming the
199
   * write buffer and to avoid false sharing on reads due to modifying the access time. The
200
   * expiration scan compensates by moving a stale-positioned head to the back of the queue when its
201
   * timestamp is fresher than the tail, so subsequent passes can resume from the next-oldest entry.
202
   *
203
   * [1] BP-Wrapper: A Framework Making Any Replacement Algorithms (Almost) Lock Contention Free
204
   * https://web.njit.edu/~dingxn/papers/BP-Wrapper.pdf
205
   * [2] TinyLFU: A Highly Efficient Cache Admission Policy
206
   * https://dl.acm.org/citation.cfm?id=3149371
207
   * [3] Adaptive Software Cache Management
208
   * https://dl.acm.org/citation.cfm?id=3274816
209
   * [4] Denial of Service via Algorithmic Complexity Attack
210
   * https://www.usenix.org/legacy/events/sec03/tech/full_papers/crosby/crosby.pdf
211
   * [5] Hashed and Hierarchical Timing Wheels
212
   * http://www.cs.columbia.edu/~nahum/w6998/papers/ton97-timing-wheels.pdf
213
   */
214

215
  static final Logger logger = System.getLogger(BoundedLocalCache.class.getName());
1 ✔
216

217
  /** The number of CPUs */
218
  static final int NCPU = Runtime.getRuntime().availableProcessors();
1 ✔
219
  /** The initial capacity of the write buffer. */
220
  static final int WRITE_BUFFER_MIN = 4;
221
  /** The maximum capacity of the write buffer. */
222
  static final int WRITE_BUFFER_MAX = 128 * ceilingPowerOfTwo(NCPU);
1 ✔
223
  /** The maximum weighted capacity of the map. */
224
  static final long MAXIMUM_CAPACITY = Long.MAX_VALUE - Integer.MAX_VALUE;
225
  /** The initial percent of the maximum weighted capacity dedicated to the main space. */
226
  static final double PERCENT_MAIN = 0.99d;
227
  /** The percent of the maximum weighted capacity dedicated to the main's protected space. */
228
  static final double PERCENT_MAIN_PROTECTED = 0.80d;
229
  /** The minimum popularity for allowing randomized admission. */
230
  static final int ADMIT_HASHDOS_THRESHOLD = 6;
231
  /** The maximum number of entries that can be transferred between queues. */
232
  static final int QUEUE_TRANSFER_THRESHOLD = 1_000;
233
  /** The maximum number of entries that can be expired per maintenance cycle. */
234
  static final int EXPIRATION_THRESHOLD = 1_000;
235
  /** The maximum number of collected references that can be drained per maintenance cycle. */
236
  static final int REFERENCE_THRESHOLD = 1_000;
237
  /** The maximum time window between touches for expiration updates. */
238
  static final long EXPIRE_TOLERANCE = TimeUnit.SECONDS.toNanos(1);
1 ✔
239
  /** The maximum duration before an entry expires. */
240
  static final long MAXIMUM_EXPIRY = (Long.MAX_VALUE >> 1); // 150 years
241
  /** The duration to wait on the eviction lock before warning of a possible misuse. */
242
  static final long WARN_AFTER_LOCK_WAIT_NANOS = TimeUnit.SECONDS.toNanos(30);
1 ✔
243
  /** The number of retries before computing to validate the entry's integrity; pow2 modulus. */
244
  static final int MAX_PUT_SPIN_WAIT_ATTEMPTS = 1024 - 1;
245
  /** The handle for the in-flight refresh operations. */
246
  static final VarHandle REFRESHES = fieldVarHandle(MethodHandles.lookup(),
1 ✔
247
      "refreshes", VarHandle.class, BoundedLocalCache.class, ConcurrentMap.class);
248

249
  final @Nullable RemovalListener<K, V> evictionListener;
250
  final @Nullable AsyncCacheLoader<K, V> cacheLoader;
251

252
  final MpscGrowableArrayQueue<Runnable> writeBuffer;
253
  final ConcurrentHashMap<Object, Node<K, V>> data;
254
  final PerformCleanupTask drainBuffersTask;
255
  final Consumer<Node<K, V>> accessPolicy;
256
  final NodeFactory<K, V> nodeFactory;
257
  final ReentrantLock evictionLock;
258
  final Weigher<K, V> weigher;
259
  final Executor executor;
260

261
  final boolean isWeighted;
262
  final boolean isAsync;
263

264
  Buffer<Node<K, V>> readBuffer;
265

266
  @Nullable Set<K> keySet;
267
  @Nullable Collection<V> values;
268
  @Nullable Set<Entry<K, V>> entrySet;
269
  volatile @Nullable ConcurrentMap<Object, CompletableFuture<?>> refreshes;
270

271
  /** Creates an instance based on the builder's configuration. */
272
  @SuppressWarnings("GuardedBy")
273
  protected BoundedLocalCache(Caffeine<K, V> builder,
274
      @Nullable AsyncCacheLoader<K, V> cacheLoader, boolean isAsync) {
1 ✔
275
    this.isAsync = isAsync;
1 ✔
276
    this.cacheLoader = cacheLoader;
1 ✔
277
    executor = builder.getExecutor();
1 ✔
278
    isWeighted = builder.isWeighted();
1 ✔
279
    evictionLock = new ReentrantLock();
1 ✔
280
    weigher = builder.getWeigher(isAsync);
1 ✔
281
    drainBuffersTask = new PerformCleanupTask(this);
1 ✔
282
    nodeFactory = NodeFactory.newFactory(builder, isAsync);
1 ✔
283
    evictionListener = builder.getEvictionListener(isAsync);
1 ✔
284
    data = new ConcurrentHashMap<>(builder.getInitialCapacity());
1 ✔
285
    writeBuffer = new MpscGrowableArrayQueue<>(WRITE_BUFFER_MIN, WRITE_BUFFER_MAX);
1 ✔
286
    boolean tracksAccess = evicts() || collectKeys() || collectValues() || expiresAfterAccess();
1 ✔
287
    readBuffer = (tracksAccess && !fastpath()) ? new BoundedBuffer<>() : Buffer.disabled();
1 ✔
288
    accessPolicy = (evicts() || expiresAfterAccess())
1 ✔
289
        ? node -> onAccess(node, Access.HIT)
1 ✔
290
        : node -> {};
1 ✔
291

292
    if (evicts()) {
1 ✔
293
      setMaximumSize(builder.getMaximum());
1 ✔
294
    }
295
  }
1 ✔
296

297
  /** Ensures that the node is alive during the map operation. */
298
  void requireIsAlive(Object key, Node<?, ?> node) {
299
    if (!node.isAlive()) {
1 ✔
300
      throw new IllegalStateException(brokenEqualityMessage(key, node));
1 ✔
301
    }
302
  }
1 ✔
303

304
  /** Logs if the node cannot be found in the map but is still alive. */
305
  void logIfAlive(Node<?, ?> node) {
306
    if (node.isAlive()) {
1 ✔
307
      String message = brokenEqualityMessage(node.getKeyReference(), node);
1 ✔
308
      logger.log(Level.ERROR, message, new IllegalStateException());
1 ✔
309
    }
310
  }
1 ✔
311

312
  /** Returns the formatted broken equality error message. */
313
  String brokenEqualityMessage(Object key, Node<?, ?> node) {
314
    return String.format(US, "An invalid state was detected, occurring when the key's equals or "
1 ✔
315
        + "hashCode was modified while residing in the cache. This violation of the Map "
316
        + "contract can lead to non-deterministic behavior (key: %s, key type: %s, "
317
        + "node type: %s, cache type: %s).", key, key.getClass().getName(),
1 ✔
318
        node.getClass().getSimpleName(), getClass().getSimpleName());
1 ✔
319
  }
320

321
  /* --------------- Shared --------------- */
322

323
  @Override
324
  public boolean isAsync() {
325
    return isAsync;
1 ✔
326
  }
327

328
  /** Returns if the node's value is currently being computed asynchronously. */
329
  final boolean isComputingAsync(@Nullable V value) {
330
    return isAsync && !Async.isReady((CompletableFuture<?>) value);
1 ✔
331
  }
332

333
  @GuardedBy("evictionLock")
334
  protected AccessOrderDeque<Node<K, V>> accessOrderWindowDeque() {
335
    throw new UnsupportedOperationException();
1 ✔
336
  }
337

338
  @GuardedBy("evictionLock")
339
  protected AccessOrderDeque<Node<K, V>> accessOrderProbationDeque() {
340
    throw new UnsupportedOperationException();
1 ✔
341
  }
342

343
  @GuardedBy("evictionLock")
344
  protected AccessOrderDeque<Node<K, V>> accessOrderProtectedDeque() {
345
    throw new UnsupportedOperationException();
1 ✔
346
  }
347

348
  @GuardedBy("evictionLock")
349
  protected WriteOrderDeque<Node<K, V>> writeOrderDeque() {
350
    throw new UnsupportedOperationException();
1 ✔
351
  }
352

353
  @Override
354
  public final Executor executor() {
355
    return executor;
1 ✔
356
  }
357

358
  @Override
359
  public ConcurrentMap<Object, CompletableFuture<?>> refreshes() {
360
    @Var var pending = refreshes;
1 ✔
361
    if (pending == null) {
1 ✔
362
      pending = new ConcurrentHashMap<>();
1 ✔
363
      if (!REFRESHES.compareAndSet(this, null, pending)) {
1 ✔
364
        pending = requireNonNull(refreshes);
1 ✔
365
      }
366
    }
367
    return pending;
1 ✔
368
  }
369

370
  /** Invalidate the in-flight refresh. */
371
  @SuppressWarnings("RedundantCollectionOperation")
372
  void discardRefresh(Object keyReference) {
373
    var pending = refreshes;
1 ✔
374
    if ((pending != null) && pending.containsKey(keyReference)) {
1 ✔
375
      pending.remove(keyReference);
1 ✔
376
    }
377
  }
1 ✔
378

379
  @Override
380
  public Object referenceKey(K key) {
381
    return nodeFactory.newLookupKey(key);
1 ✔
382
  }
383

384
  @Override
385
  public boolean isPendingEviction(K key) {
386
    Node<K, V> node = data.get(nodeFactory.newLookupKey(key));
1 ✔
387
    if (node == null) {
1 ✔
388
      return false;
1 ✔
389
    }
390
    boolean expired = hasExpired(node, expirationTicker().read());
1 ✔
391
    V value = node.getValue();
1 ✔
392
    return (value == null) || (expired && !isComputingAsync(value));
1 !
393
  }
394

395
  /* --------------- Stats Support --------------- */
396

397
  @Override
398
  public boolean isRecordingStats() {
399
    return false;
1 ✔
400
  }
401

402
  @Override
403
  public StatsCounter statsCounter() {
404
    return StatsCounter.disabledStatsCounter();
1 ✔
405
  }
406

407
  @Override
408
  public Ticker statsTicker() {
409
    return Ticker.disabledTicker();
1 ✔
410
  }
411

412
  /* --------------- Removal Listener Support --------------- */
413

414
  @Override
415
  public @Nullable RemovalListener<K, V> removalListener() {
416
    return null;
1 ✔
417
  }
418

419
  @Override
420
  public void notifyRemoval(@Nullable K key, @Nullable V value, RemovalCause cause) {
421
    var removalListener = removalListener();
1 ✔
422
    if (removalListener == null) {
1 ✔
423
      return;
1 ✔
424
    }
425
    Runnable task = () -> {
1 ✔
426
      try {
427
        removalListener.onRemoval(key, value, cause);
1 ✔
428
      } catch (Throwable t) {
1 ✔
429
        logger.log(Level.WARNING, "Exception thrown by removal listener", t);
1 ✔
430
      }
1 ✔
431
    };
1 ✔
432
    try {
433
      executor.execute(task);
1 ✔
434
    } catch (Throwable t) {
1 ✔
435
      logger.log(Level.ERROR, "Exception thrown when submitting removal listener", t);
1 ✔
436
      task.run();
1 ✔
437
    }
1 ✔
438
  }
1 ✔
439

440
  /* --------------- Eviction Listener Support --------------- */
441

442
  void notifyEviction(@Nullable K key, @Nullable V value, RemovalCause cause) {
443
    if (evictionListener == null) {
1 ✔
444
      return;
1 ✔
445
    }
446
    try {
447
      evictionListener.onRemoval(key, value, cause);
1 ✔
448
    } catch (Throwable t) {
1 ✔
449
      logger.log(Level.WARNING, "Exception thrown by eviction listener", t);
1 ✔
450
    }
1 ✔
451
  }
1 ✔
452

453
  /* --------------- Reference Support --------------- */
454

455
  @Override
456
  public boolean collectKeys() {
457
    return false;
1 ✔
458
  }
459

460
  /** Returns if the values are weak or soft reference garbage collected. */
461
  protected boolean collectValues() {
462
    return false;
1 ✔
463
  }
464

465
  @SuppressWarnings({"DataFlowIssue", "NullAway"})
466
  protected ReferenceQueue<K> keyReferenceQueue() {
467
    return null;
1 ✔
468
  }
469

470
  @SuppressWarnings({"DataFlowIssue", "NullAway"})
471
  protected ReferenceQueue<V> valueReferenceQueue() {
472
    return null;
1 ✔
473
  }
474

475
  /* --------------- Expiration Support --------------- */
476

477
  /** Returns the {@link Pacer} used to schedule the maintenance task. */
478
  protected @Nullable Pacer pacer() {
479
    return null;
1 ✔
480
  }
481

482
  /** Returns if the cache expires entries by any of its policies. */
483
  final boolean expires() {
484
    return expiresAfterAccess() || expiresAfterWrite() || expiresVariable();
1 ✔
485
  }
486

487
  /** Returns if a read may extend when the entry expires. */
488
  final boolean expiresAfterRead() {
489
    return expiresAfterAccess() || expiresVariable();
1 ✔
490
  }
491

492
  /** Returns if the cache expires entries after a variable time threshold. */
493
  protected boolean expiresVariable() {
494
    return false;
1 ✔
495
  }
496

497
  /** Returns if the cache expires entries after an access time threshold. */
498
  protected boolean expiresAfterAccess() {
499
    return false;
1 ✔
500
  }
501

502
  /** Returns how long after the last access to an entry the map will retain that entry. */
503
  protected long expiresAfterAccessNanos() {
504
    throw new UnsupportedOperationException();
1 ✔
505
  }
506

507
  protected void setExpiresAfterAccessNanos(long expireAfterAccessNanos) {
508
    throw new UnsupportedOperationException();
1 ✔
509
  }
510

511
  /** Returns if the cache expires entries after a write time threshold. */
512
  protected boolean expiresAfterWrite() {
513
    return false;
1 ✔
514
  }
515

516
  /** Returns how long after the last write to an entry the map will retain that entry. */
517
  protected long expiresAfterWriteNanos() {
518
    throw new UnsupportedOperationException();
1 ✔
519
  }
520

521
  protected void setExpiresAfterWriteNanos(long expireAfterWriteNanos) {
522
    throw new UnsupportedOperationException();
1 ✔
523
  }
524

525
  /** Returns if the cache refreshes entries after a write time threshold. */
526
  protected boolean refreshAfterWrite() {
527
    return false;
1 ✔
528
  }
529

530
  /** Returns how long after the last write an entry becomes a candidate for refresh. */
531
  protected long refreshAfterWriteNanos() {
532
    throw new UnsupportedOperationException();
1 ✔
533
  }
534

535
  protected void setRefreshAfterWriteNanos(long refreshAfterWriteNanos) {
536
    throw new UnsupportedOperationException();
1 ✔
537
  }
538

539
  @Override
540
  @SuppressWarnings({"DataFlowIssue", "NullAway"})
541
  public Expiry<K, V> expiry() {
542
    return null;
1 ✔
543
  }
544

545
  /** Returns the {@link Ticker} used by this cache for expiration. */
546
  public Ticker expirationTicker() {
547
    return Ticker.disabledTicker();
1 ✔
548
  }
549

550
  protected TimerWheel<K, V> timerWheel() {
551
    throw new UnsupportedOperationException();
1 ✔
552
  }
553

554
  /* --------------- Eviction Support --------------- */
555

556
  /** Returns if the cache evicts entries due to a maximum size or weight threshold. */
557
  protected boolean evicts() {
558
    return false;
1 ✔
559
  }
560

561
  /** Returns if entries may be assigned different weights. */
562
  protected boolean isWeighted() {
563
    return isWeighted && (weigher != Weigher.singletonWeigher());
1 ✔
564
  }
565

566
  protected FrequencySketch frequencySketch() {
567
    throw new UnsupportedOperationException();
1 ✔
568
  }
569

570
  @GuardedBy("evictionLock")
571
  protected WindowClimber climber() {
572
    throw new UnsupportedOperationException();
1 ✔
573
  }
574

575
  /** Returns if an access to an entry can skip notifying the eviction policy. */
576
  protected boolean fastpath() {
577
    return false;
1 ✔
578
  }
579

580
  /** Returns the maximum weighted size. */
581
  protected long maximum() {
582
    throw new UnsupportedOperationException();
1 ✔
583
  }
584

585
  /** Returns the maximum weighted size. */
586
  protected long maximumAcquire() {
587
    throw new UnsupportedOperationException();
1 ✔
588
  }
589

590
  /** Returns the maximum weighted size of the window space. */
591
  protected long windowMaximum() {
592
    throw new UnsupportedOperationException();
1 ✔
593
  }
594

595
  /** Returns the maximum weighted size of the main's protected space. */
596
  protected long mainProtectedMaximum() {
597
    throw new UnsupportedOperationException();
1 ✔
598
  }
599

600
  @GuardedBy("evictionLock")
601
  protected void setMaximum(long maximum) {
602
    throw new UnsupportedOperationException();
1 ✔
603
  }
604

605
  @GuardedBy("evictionLock")
606
  protected void setWindowMaximum(long maximum) {
607
    throw new UnsupportedOperationException();
1 ✔
608
  }
609

610
  @GuardedBy("evictionLock")
611
  protected void setMainProtectedMaximum(long maximum) {
612
    throw new UnsupportedOperationException();
1 ✔
613
  }
614

615
  /** Returns the combined weight of the values in the cache (may be negative). */
616
  protected long weightedSize() {
617
    throw new UnsupportedOperationException();
1 ✔
618
  }
619

620
  /** Returns the combined weight of the values in the cache (may be negative). */
621
  protected long weightedSizeAcquire() {
622
    throw new UnsupportedOperationException();
1 ✔
623
  }
624

625
  /** Returns the uncorrected combined weight of the values in the window space. */
626
  protected long windowWeightedSize() {
627
    throw new UnsupportedOperationException();
1 ✔
628
  }
629

630
  /** Returns the uncorrected combined weight of the values in the main's protected space. */
631
  protected long mainProtectedWeightedSize() {
632
    throw new UnsupportedOperationException();
1 ✔
633
  }
634

635
  @GuardedBy("evictionLock")
636
  protected void setWeightedSize(long weightedSize) {
637
    throw new UnsupportedOperationException();
1 ✔
638
  }
639

640
  @GuardedBy("evictionLock")
641
  protected void setWindowWeightedSize(long weightedSize) {
642
    throw new UnsupportedOperationException();
1 ✔
643
  }
644

645
  @GuardedBy("evictionLock")
646
  protected void setMainProtectedWeightedSize(long weightedSize) {
647
    throw new UnsupportedOperationException();
1 ✔
648
  }
649

650
  /**
651
   * Sets the maximum weighted size of the cache. The caller may need to perform a maintenance cycle
652
   * to eagerly evicts entries until the cache shrinks to the appropriate size.
653
   */
654
  @GuardedBy("evictionLock")
655
  @SuppressWarnings({"ConstantValue", "Varifier"})
656
  void setMaximumSize(long maximum) {
657
    requireArgument(maximum >= 0, "maximum must not be negative");
1 ✔
658
    long max = Math.min(maximum, MAXIMUM_CAPACITY);
1 ✔
659
    if (max == maximum()) {
1 ✔
660
      return;
1 ✔
661
    }
662

663
    long window = max - (long) (PERCENT_MAIN * max);
1 ✔
664
    long mainProtected = (long) (PERCENT_MAIN_PROTECTED * (max - window));
1 ✔
665

666
    setMaximum(max);
1 ✔
667
    setWindowMaximum(window);
1 ✔
668
    setMainProtectedMaximum(mainProtected);
1 ✔
669

670
    if (climber() != null) {
1 ✔
671
      // null during the super constructor's initial sizing; the generated constructor replays this
672
      climber().resized(maximum());
1 ✔
673
    }
674

675
    if ((frequencySketch() != null) && !isWeighted() && (weightedSize() >= (max >>> 1))) {
1 ✔
676
      // Lazily initialize when close to the maximum size
677
      frequencySketch().ensureCapacity(max);
1 ✔
678
      recordReads();
1 ✔
679
    }
680
  }
1 ✔
681

682
  /** Evicts entries if the cache exceeds the maximum. */
683
  @GuardedBy("evictionLock")
684
  void evictEntries(long now) {
685
    if (!evicts()) {
1 ✔
686
      return;
1 ✔
687
    }
688
    var candidate = evictFromWindow();
1 ✔
689
    evictFromMain(candidate, now);
1 ✔
690
  }
1 ✔
691

692
  /**
693
   * Evicts entries from the window space into the main space while the window size exceeds a
694
   * maximum, transferring at most a bounded number of entries per maintenance cycle.
695
   *
696
   * @return the first candidate promoted into the probation space
697
   */
698
  @GuardedBy("evictionLock")
699
  @Nullable Node<K, V> evictFromWindow() {
700
    @Var Node<K, V> first = null;
1 ✔
701
    @Var int remaining = QUEUE_TRANSFER_THRESHOLD;
1 ✔
702
    @Var Node<K, V> node = accessOrderWindowDeque().peekFirst();
1 ✔
703
    while (windowWeightedSize() > windowMaximum()) {
1 ✔
704
      // The pending operations will adjust the size to reflect the correct weight
705
      if (node == null) {
1 ✔
706
        break;
1 ✔
707
      }
708
      if (remaining == 0) {
1 ✔
709
        // A resize can leave the window arbitrarily oversized, so the transfer is bounded per
710
        // maintenance cycle and re-armed to drain the backlog across cycles
711
        setDrainStatusOpaque(PROCESSING_TO_REQUIRED);
1 ✔
712
        break;
1 ✔
713
      }
714

715
      Node<K, V> next = node.getNextInAccessOrder();
1 ✔
716
      long weight = node.getPolicyWeight();
1 ✔
717
      if (weight != 0) {
1 ✔
718
        transfer(node, weight, WINDOW, PROBATION);
1 ✔
719
        if (first == null) {
1 ✔
720
          first = node;
1 ✔
721
        }
722
        remaining--;
1 ✔
723
      }
724
      node = next;
1 ✔
725
    }
1 ✔
726

727
    return first;
1 ✔
728
  }
729

730
  /**
731
   * Evicts entries from the main space if the cache exceeds the maximum capacity. The main space
732
   * determines whether admitting an entry (coming from the window space) is preferable to retaining
733
   * the eviction policy's victim. This decision is made using a frequency filter so that the
734
   * least frequently used entry is removed.
735
   * <p>
736
   * The window space's candidates were previously promoted to the probation space at its MRU
737
   * position and the eviction policy's victim starts at the LRU position. The candidates are
738
   * evaluated in promotion order while an eviction is required, and if exhausted then additional
739
   * entries are retrieved from the window space. Likewise, if the victim selection exhausts the
740
   * probation space then additional entries are retrieved from the protected space. The queues are
741
   * consumed in LRU order and the evicted entry is the one with a lower relative frequency, where
742
   * the preference is to retain the main space's victims versus the window space's candidates on a
743
   * tie.
744
   *
745
   * @param candidate the first candidate promoted into the probation space
746
   * @param now the current time, in nanoseconds
747
   */
748
  @GuardedBy("evictionLock")
749
  void evictFromMain(@Var @Nullable Node<K, V> candidate, long now) {
750
    @Var int victimQueue = PROBATION;
1 ✔
751
    @Var int candidateQueue = PROBATION;
1 ✔
752
    @Var Node<K, V> victim = accessOrderProbationDeque().peekFirst();
1 ✔
753
    while (weightedSize() > maximum()) {
1 ✔
754
      // Search the admission window for additional candidates
755
      if ((candidate == null) && (candidateQueue == PROBATION)) {
1 ✔
756
        candidate = accessOrderWindowDeque().peekFirst();
1 ✔
757
        candidateQueue = WINDOW;
1 ✔
758
      }
759

760
      // Try evicting from the protected and window queues
761
      if ((candidate == null) && (victim == null)) {
1 ✔
762
        if (victimQueue == PROBATION) {
1 ✔
763
          victim = accessOrderProtectedDeque().peekFirst();
1 ✔
764
          victimQueue = PROTECTED;
1 ✔
765
          continue;
1 ✔
766
        } else if (victimQueue == PROTECTED) {
1 ✔
767
          victim = accessOrderWindowDeque().peekFirst();
1 ✔
768
          victimQueue = WINDOW;
1 ✔
769
          continue;
1 ✔
770
        }
771

772
        // The pending operations will adjust the size to reflect the correct weight
773
        break;
774
      }
775

776
      // Skip over entries with zero weight
777
      if ((victim != null) && (victim.getPolicyWeight() == 0)) {
1 ✔
778
        victim = victim.getNextInAccessOrder();
1 ✔
779
        continue;
1 ✔
780
      } else if ((candidate != null) && (candidate.getPolicyWeight() == 0)) {
1 ✔
781
        candidate = candidate.getNextInAccessOrder();
1 ✔
782
        continue;
1 ✔
783
      }
784

785
      // Evict immediately if only one of the entries is present
786
      if (victim == null) {
1 ✔
787
        requireNonNull(candidate);
1 ✔
788
        candidate = evictAndAdvance(candidate, now, RemovalCause.SIZE);
1 ✔
789
        continue;
1 ✔
790
      } else if (candidate == null) {
1 ✔
791
        victim = evictAndAdvance(victim, now, RemovalCause.SIZE);
1 ✔
792
        continue;
1 ✔
793
      }
794

795
      // Evict immediately if both selected the same entry
796
      if (candidate == victim) {
1 ✔
797
        victim = evictAndAdvance(victim, now, RemovalCause.SIZE);
1 ✔
798
        candidate = null;
1 ✔
799
        continue;
1 ✔
800
      }
801

802
      // Evict immediately if an entry was collected
803
      var victimKeyRef = victim.getKeyReferenceOrNull();
1 ✔
804
      var candidateKeyRef = candidate.getKeyReferenceOrNull();
1 ✔
805
      if (victimKeyRef == null) {
1 ✔
806
        victim = evictAndAdvance(victim, now, RemovalCause.COLLECTED);
1 ✔
807
        continue;
1 ✔
808
      } else if (candidateKeyRef == null) {
1 ✔
809
        candidate = evictAndAdvance(candidate, now, RemovalCause.COLLECTED);
1 ✔
810
        continue;
1 ✔
811
      }
812

813
      // Evict immediately if an entry was removed
814
      if (!victim.isAlive()) {
1 ✔
815
        victim = evictAndAdvance(victim, now, RemovalCause.SIZE);
1 ✔
816
        continue;
1 ✔
817
      } else if (!candidate.isAlive()) {
1 ✔
818
        candidate = evictAndAdvance(candidate, now, RemovalCause.SIZE);
1 ✔
819
        continue;
1 ✔
820
      }
821

822
      // Evict immediately if the candidate's weight exceeds the maximum
823
      if (candidate.getPolicyWeight() > maximum()) {
1 ✔
824
        candidate = evictAndAdvance(candidate, now, RemovalCause.SIZE);
1 ✔
825
        continue;
1 ✔
826
      }
827

828
      // Evict the entry with the lowest frequency
829
      if (admit(candidateKeyRef, victimKeyRef)) {
1 ✔
830
        victim = evictAndAdvance(victim, now, RemovalCause.SIZE);
1 ✔
831
        candidate = candidate.getNextInAccessOrder();
1 ✔
832
      } else {
833
        candidate = evictAndAdvance(candidate, now, RemovalCause.SIZE);
1 ✔
834
      }
835
    }
1 ✔
836
  }
1 ✔
837

838
  /** Evicts the entry and returns its successor, or {@code null} if it was the last. */
839
  @GuardedBy("evictionLock")
840
  @Nullable Node<K, V> evictAndAdvance(Node<K, V> node, long now, RemovalCause cause) {
841
    Node<K, V> next = node.getNextInAccessOrder();
1 ✔
842
    evictEntry(node, cause, now);
1 ✔
843
    return next;
1 ✔
844
  }
845

846
  /**
847
   * Determines if the candidate should be accepted into the main space, as determined by its
848
   * frequency relative to the victim. A small amount of randomness is used to protect against hash
849
   * collision attacks, where the victim's frequency is artificially raised so that no new entries
850
   * are admitted.
851
   *
852
   * @param candidateKeyRef the keyRef for the entry being proposed for long term retention
853
   * @param victimKeyRef the keyRef for the entry chosen by the eviction policy for replacement
854
   * @return if the candidate should be admitted and the victim ejected
855
   */
856
  @GuardedBy("evictionLock")
857
  boolean admit(Object candidateKeyRef, Object victimKeyRef) {
858
    int candidateFreq = frequencySketch().frequency(candidateKeyRef);
1 ✔
859
    int victimFreq = frequencySketch().frequency(victimKeyRef);
1 ✔
860
    if (candidateFreq > victimFreq) {
1 ✔
861
      return true;
1 ✔
862
    } else if (candidateFreq >= ADMIT_HASHDOS_THRESHOLD) {
1 ✔
863
      // The maximum frequency is 15 and halved to 7 after a reset to age the history. An attack
864
      // exploits that a hot candidate is rejected in favor of a hot victim. The threshold of a warm
865
      // candidate reduces the number of random acceptances to minimize the impact on the hit rate.
866
      int random = ThreadLocalRandom.current().nextInt();
1 ✔
867
      return ((random & 127) == 0);
1 ✔
868
    }
869
    return false;
1 ✔
870
  }
871

872
  /** Expires entries that have expired by access, write, or variable. */
873
  @GuardedBy("evictionLock")
874
  void expireEntries(long now) {
875
    expireAfterAccessEntries(now);
1 ✔
876
    expireAfterWriteEntries(now);
1 ✔
877
    expireVariableEntries(now);
1 ✔
878

879
    Pacer pacer = pacer();
1 ✔
880
    if (pacer != null) {
1 ✔
881
      long delay = getExpirationDelay(now);
1 ✔
882
      if (delay == Long.MAX_VALUE) {
1 ✔
883
        pacer.cancel();
1 ✔
884
      } else {
885
        pacer.schedule(executor, drainBuffersTask, now, delay);
1 ✔
886
      }
887
    }
888
  }
1 ✔
889

890
  /** Expires entries in the access-order queue. */
891
  @GuardedBy("evictionLock")
892
  void expireAfterAccessEntries(long now) {
893
    if (!expiresAfterAccess()) {
1 ✔
894
      return;
1 ✔
895
    }
896

897
    @Var int remaining = EXPIRATION_THRESHOLD;
1 ✔
898
    remaining = expireAfterAccessEntries(now, WINDOW, accessOrderWindowDeque(), remaining);
1 ✔
899
    if (evicts()) {
1 ✔
900
      remaining = expireAfterAccessEntries(
1 ✔
901
          now, PROBATION, accessOrderProbationDeque(), remaining);
1 ✔
902
      remaining = expireAfterAccessEntries(
1 ✔
903
          now, PROTECTED, accessOrderProtectedDeque(), remaining);
1 ✔
904
    }
905
    if (remaining == 0) {
1 ✔
906
      setDrainStatusOpaque(PROCESSING_TO_REQUIRED);
1 ✔
907
    }
908
  }
1 ✔
909

910
  /**
911
   * Expires entries in an access-order queue, up to the {@code remaining} budget, and returns the
912
   * unused budget. When exhausted the caller re-arms maintenance to process the backlog.
913
   */
914
  @GuardedBy("evictionLock")
915
  int expireAfterAccessEntries(long now, int queueType,
916
      AccessOrderDeque<Node<K, V>> accessOrderDeque, @Var int remaining) {
917
    var head = accessOrderDeque.peekFirst();
1 ✔
918
    if (head == null) {
1 ✔
919
      return remaining;
1 ✔
920
    }
921
    long duration = expiresAfterAccessNanos();
1 ✔
922
    var last = requireNonNull(accessOrderDeque.peekLast());
1 ✔
923
    for (var node = head; (node != null) && (remaining > 0);) {
1 ✔
924
      var next = (node == last) ? null : node.getNextInAccessOrder();
1 ✔
925
      if ((now - node.getAccessTime()) < duration) {
1 ✔
926
        boolean stalePosition = ((last.getAccessTime() - node.getAccessTime()) < 0);
1 ✔
927
        if (stalePosition || isComputingAsync(node.getValue())) {
1 ✔
928
          reorder(accessOrderDeque, node, queueType);
1 ✔
929
          node = next;
1 ✔
930
          continue;
1 ✔
931
        }
932
        return remaining;
1 ✔
933
      }
934
      evictEntry(node, RemovalCause.EXPIRED, now);
1 ✔
935
      remaining--;
1 ✔
936
      node = next;
1 ✔
937

938
      // Reentrant maintenance can remove the captured tail, preventing termination
939
      boolean bounded = (last.getQueueType() == queueType) && accessOrderDeque.contains(last);
1 ✔
940
      if ((node != null) && !bounded) {
1 ✔
941
        return 0;
1 ✔
942
      }
943
    }
1 ✔
944
    return remaining;
1 ✔
945
  }
946

947
  /** Expires entries on the write-order queue. */
948
  @GuardedBy("evictionLock")
949
  void expireAfterWriteEntries(long now) {
950
    if (!expiresAfterWrite()) {
1 ✔
951
      return;
1 ✔
952
    }
953

954
    var head = writeOrderDeque().peekFirst();
1 ✔
955
    if (head == null) {
1 ✔
956
      return;
1 ✔
957
    }
958
    long duration = expiresAfterWriteNanos();
1 ✔
959
    @Var int remaining = EXPIRATION_THRESHOLD;
1 ✔
960
    var last = requireNonNull(writeOrderDeque().peekLast());
1 ✔
961
    for (var node = head; (node != null) && (remaining > 0);) {
1 ✔
962
      var next = (node == last) ? null : node.getNextInWriteOrder();
1 ✔
963
      if ((now - node.getWriteTime()) < duration) {
1 ✔
964
        boolean stalePosition = ((last.getWriteTime() - node.getWriteTime()) < 0);
1 ✔
965
        if (stalePosition || isComputingAsync(node.getValue())) {
1 ✔
966
          reorder(writeOrderDeque(), node);
1 ✔
967
          node = next;
1 ✔
968
          continue;
1 ✔
969
        }
970
        return;
1 ✔
971
      }
972
      evictEntry(node, RemovalCause.EXPIRED, now);
1 ✔
973
      remaining--;
1 ✔
974
      node = next;
1 ✔
975

976
      // Reentrant maintenance can remove the captured tail, preventing termination
977
      if ((node != null) && !writeOrderDeque().contains(last)) {
1 ✔
978
        remaining = 0;
1 ✔
979
        break;
1 ✔
980
      }
981
    }
1 ✔
982
    if (remaining == 0) {
1 ✔
983
      setDrainStatusOpaque(PROCESSING_TO_REQUIRED);
1 ✔
984
    }
985
  }
1 ✔
986

987
  /** Expires entries in the timer wheel. */
988
  @GuardedBy("evictionLock")
989
  void expireVariableEntries(long now) {
990
    if (expiresVariable() && (timerWheel().advance(this, now, EXPIRATION_THRESHOLD) == 0)) {
1 ✔
991
      setDrainStatusOpaque(PROCESSING_TO_REQUIRED);
1 ✔
992
    }
993
  }
1 ✔
994

995
  /** Returns the duration until the next item expires, or {@link Long#MAX_VALUE} if none. */
996
  @GuardedBy("evictionLock")
997
  long getExpirationDelay(long now) {
998
    @Var long delay = Long.MAX_VALUE;
1 ✔
999
    if (expiresAfterAccess()) {
1 ✔
1000
      @Var Node<K, V> node = accessOrderWindowDeque().peekFirst();
1 ✔
1001
      if (node != null) {
1 ✔
1002
        long age = Math.max(0, now - node.getAccessTime());
1 ✔
1003
        delay = Math.min(delay, expiresAfterAccessNanos() - age);
1 ✔
1004
      }
1005
      if (evicts()) {
1 ✔
1006
        node = accessOrderProbationDeque().peekFirst();
1 ✔
1007
        if (node != null) {
1 ✔
1008
          long age = Math.max(0, now - node.getAccessTime());
1 ✔
1009
          delay = Math.min(delay, expiresAfterAccessNanos() - age);
1 ✔
1010
        }
1011
        node = accessOrderProtectedDeque().peekFirst();
1 ✔
1012
        if (node != null) {
1 ✔
1013
          long age = Math.max(0, now - node.getAccessTime());
1 ✔
1014
          delay = Math.min(delay, expiresAfterAccessNanos() - age);
1 ✔
1015
        }
1016
      }
1017
    }
1018
    if (expiresAfterWrite()) {
1 ✔
1019
      Node<K, V> node = writeOrderDeque().peekFirst();
1 ✔
1020
      if (node != null) {
1 ✔
1021
        long age = Math.max(0, now - node.getWriteTime());
1 ✔
1022
        delay = Math.min(delay, expiresAfterWriteNanos() - age);
1 ✔
1023
      }
1024
    }
1025
    if (expiresVariable()) {
1 ✔
1026
      delay = Math.min(delay, timerWheel().getExpirationDelay());
1 ✔
1027
    }
1028
    return delay;
1 ✔
1029
  }
1030

1031
  /** Returns if the entry has expired. The caller must exempt an in-flight asynchronous load. */
1032
  @SuppressWarnings("ShortCircuitBoolean")
1033
  boolean hasExpired(Node<K, V> node, long now) {
1034
    if (!expires()) {
1 ✔
1035
      return false;
1 ✔
1036
    }
1037
    boolean expired =
1 ✔
1038
        (expiresAfterAccess() && (now - node.getAccessTime() >= expiresAfterAccessNanos()))
1 ✔
1039
        | (expiresAfterWrite() && (now - node.getWriteTime() >= expiresAfterWriteNanos()))
1 ✔
1040
        | (expiresVariable() && (now - node.getVariableTime() >= 0));
1 ✔
1041
    VarHandle.loadLoadFence();
1 ✔
1042
    return expired;
1 ✔
1043
  }
1044

1045
  /**
1046
   * Attempts to evict the entry based on the given removal cause. A removal may be ignored if the
1047
   * entry was updated and is no longer eligible for eviction.
1048
   *
1049
   * @param node the entry to evict
1050
   * @param cause the reason to evict
1051
   * @param now the current time, used only if expiring
1052
   * @return if the entry was evicted
1053
   */
1054
  @GuardedBy("evictionLock")
1055
  @SuppressWarnings({"GuardedByChecker", "SynchronizationOnLocalVariableOrMethodParameter"})
1056
  boolean evictEntry(Node<K, V> node, RemovalCause cause, long now) {
1057
    K key = node.getKey();
1 ✔
1058
    var ctx = new EvictContext<V>();
1 ✔
1059
    var keyReference = node.getKeyReference();
1 ✔
1060

1061
    data.computeIfPresent(keyReference, (k, n) -> {
1 ✔
1062
      if (n != node) {
1 !
UNCOV
1063
        return n;
×
1064
      }
1065
      synchronized (node) {
1 ✔
1066
        ctx.value = node.getValue();
1 ✔
1067
        boolean expired = hasExpired(node, now);
1 ✔
1068
        boolean computing = isComputingAsync(ctx.value);
1 ✔
1069
        if ((key == null) || (ctx.value == null)) {
1 ✔
1070
          ctx.cause = RemovalCause.COLLECTED;
1 ✔
1071
        } else if (cause == RemovalCause.COLLECTED) {
1 ✔
1072
          ctx.resurrect = true;
1 ✔
1073
          return node;
1 ✔
1074
        } else if (expired && !computing) {
1 ✔
1075
          ctx.cause = RemovalCause.EXPIRED;
1 ✔
1076
        } else {
1077
          ctx.cause = cause;
1 ✔
1078
        }
1079

1080
        if (ctx.cause == RemovalCause.EXPIRED) {
1 ✔
1081
          if (!expired) {
1 ✔
1082
            ctx.resurrect = true;
1 ✔
1083
            return node;
1 ✔
1084
          } else if (computing) {
1 ✔
1085
            long sentinel = (now + ASYNC_EXPIRY);
1 ✔
1086
            setVariableTime(node, sentinel);
1 ✔
1087
            setAccessTime(node, sentinel);
1 ✔
1088
            setWriteTime(node, sentinel);
1 ✔
1089
            ctx.resurrect = true;
1 ✔
1090
            return node;
1 ✔
1091
          }
1092
        } else if (ctx.cause == RemovalCause.SIZE) {
1 ✔
1093
          int weight = node.getWeight();
1 ✔
1094
          if (weight == 0) {
1 ✔
1095
            ctx.resurrect = true;
1 ✔
1096
            return node;
1 ✔
1097
          }
1098
        }
1099

1100
        notifyEviction(key, ctx.value, ctx.cause);
1 ✔
1101
        discardRefresh(keyReference);
1 ✔
1102
        ctx.removed = true;
1 ✔
1103
        node.retire();
1 ✔
1104
        return null;
1 ✔
1105
      }
1106
    });
1107

1108
    // The entry is no longer eligible for eviction
1109
    if (ctx.resurrect) {
1 ✔
1110
      return false;
1 ✔
1111
    }
1112

1113
    // If the eviction fails due to a concurrent removal of the victim, that removal may cancel out
1114
    // the addition that triggered this eviction. The victim is eagerly unlinked and the size
1115
    // decremented before the removal task so that if an eviction is still required then a new
1116
    // victim will be chosen for removal.
1117
    unlink(node);
1 ✔
1118

1119
    synchronized (node) {
1 ✔
1120
      logIfAlive(node);
1 ✔
1121
      makeDead(node);
1 ✔
1122
    }
1 ✔
1123

1124
    if (ctx.removed) {
1 ✔
1125
      var removeCause = requireNonNull(ctx.cause);
1 ✔
1126
      statsCounter().recordEviction(node.getWeight(), removeCause);
1 ✔
1127
      notifyRemoval(key, ctx.value, removeCause);
1 ✔
1128
    }
1129

1130
    return true;
1 ✔
1131
  }
1132

1133
  /** Adapts the eviction policy to towards the optimal recency / frequency configuration. */
1134
  @GuardedBy("evictionLock")
1135
  @SuppressWarnings("UnnecessaryReturnStatement")
1136
  void climb() {
1137
    if (!evicts()) {
1 ✔
1138
      return;
1 ✔
1139
    }
1140

1141
    if (frequencySketch().isNotInitialized()) {
1 ✔
1142
      climber().resetSample();
1 ✔
1143
    } else if (weightedSize() >= (maximum() >>> 1)) {
1 ✔
1144
      climber().determineAdjustment(maximum(), windowMaximum(),
1 ✔
1145
          mainProtectedMaximum(), frequencySketch().sampleSize);
1 ✔
1146
    } else {
1147
      climber().discardSample(maximum(), frequencySketch().sampleSize);
1 ✔
1148
    }
1149

1150
    demoteFromMainProtected();
1 ✔
1151
    long amount = climber().adjustment();
1 ✔
1152
    if (amount == 0) {
1 ✔
1153
      return;
1 ✔
1154
    } else if (amount > 0) {
1 ✔
1155
      increaseWindow();
1 ✔
1156
    } else {
1157
      decreaseWindow();
1 ✔
1158
    }
1159
  }
1 ✔
1160

1161
  /**
1162
   * Increases the size of the admission window by shrinking the portion allocated to the main
1163
   * space. As the main space is partitioned into probation and protected regions (80% / 20%), for
1164
   * simplicity only the protected is reduced. If the regions exceed their maximums, this may cause
1165
   * protected items to be demoted to the probation region and probation items to be demoted to the
1166
   * admission window.
1167
   */
1168
  @GuardedBy("evictionLock")
1169
  void increaseWindow() {
1170
    if (mainProtectedMaximum() == 0) {
1 ✔
1171
      return;
1 ✔
1172
    }
1173

1174
    @SuppressWarnings("MathClampLong")
1175
    @Var long quota = Math.min(climber().adjustment(), mainProtectedMaximum());
1 ✔
1176
    setMainProtectedMaximum(mainProtectedMaximum() - quota);
1 ✔
1177
    setWindowMaximum(windowMaximum() + quota);
1 ✔
1178
    demoteFromMainProtected();
1 ✔
1179

1180
    for (int i = 0; i < QUEUE_TRANSFER_THRESHOLD; i++) {
1 ✔
1181
      @Var Node<K, V> candidate = accessOrderProbationDeque().peekFirst();
1 ✔
1182
      @Var boolean probation = true;
1 ✔
1183
      if ((candidate == null) || (quota < candidate.getPolicyWeight())) {
1 ✔
1184
        candidate = accessOrderProtectedDeque().peekFirst();
1 ✔
1185
        probation = false;
1 ✔
1186
      }
1187
      if (candidate == null) {
1 ✔
1188
        break;
1 ✔
1189
      }
1190

1191
      long weight = candidate.getPolicyWeight();
1 ✔
1192
      if (quota < weight) {
1 ✔
1193
        break;
1 ✔
1194
      }
1195

1196
      quota -= weight;
1 ✔
1197
      transfer(candidate, weight, probation ? PROBATION : PROTECTED, WINDOW);
1 ✔
1198
    }
1199

1200
    setMainProtectedMaximum(mainProtectedMaximum() + quota);
1 ✔
1201
    setWindowMaximum(windowMaximum() - quota);
1 ✔
1202
    climber().carryOver(quota);
1 ✔
1203
  }
1 ✔
1204

1205
  /** Decreases the size of the admission window and increases the main's protected region. */
1206
  @GuardedBy("evictionLock")
1207
  void decreaseWindow() {
1208
    if (windowMaximum() <= 1) {
1 ✔
1209
      return;
1 ✔
1210
    }
1211

1212
    @SuppressWarnings("MathClampLong")
1213
    @Var long quota = Math.min(-climber().adjustment(), Math.max(0, windowMaximum() - 1));
1 ✔
1214
    setMainProtectedMaximum(mainProtectedMaximum() + quota);
1 ✔
1215
    setWindowMaximum(windowMaximum() - quota);
1 ✔
1216

1217
    for (int i = 0; i < QUEUE_TRANSFER_THRESHOLD; i++) {
1 ✔
1218
      Node<K, V> candidate = accessOrderWindowDeque().peekFirst();
1 ✔
1219
      if (candidate == null) {
1 ✔
1220
        break;
1 ✔
1221
      }
1222

1223
      long weight = candidate.getPolicyWeight();
1 ✔
1224
      if (quota < weight) {
1 ✔
1225
        break;
1 ✔
1226
      }
1227

1228
      quota -= weight;
1 ✔
1229
      transfer(candidate, weight, WINDOW, PROBATION);
1 ✔
1230
    }
1231

1232
    setMainProtectedMaximum(mainProtectedMaximum() - quota);
1 ✔
1233
    setWindowMaximum(windowMaximum() + quota);
1 ✔
1234
    climber().carryOver(-quota);
1 ✔
1235
  }
1 ✔
1236

1237
  /** Transfers the nodes from the protected to the probation region if it exceeds the maximum. */
1238
  @GuardedBy("evictionLock")
1239
  void demoteFromMainProtected() {
1240
    long mainProtectedMaximum = mainProtectedMaximum();
1 ✔
1241
    @Var int remaining = QUEUE_TRANSFER_THRESHOLD;
1 ✔
1242
    while (mainProtectedWeightedSize() > mainProtectedMaximum) {
1 ✔
1243
      Node<K, V> demoted = accessOrderProtectedDeque().peekFirst();
1 ✔
1244
      if (demoted == null) {
1 ✔
1245
        break;
1 ✔
1246
      }
1247
      if (remaining == 0) {
1 ✔
1248
        setDrainStatusOpaque(PROCESSING_TO_REQUIRED);
1 ✔
1249
        break;
1 ✔
1250
      }
1251
      transfer(demoted, demoted.getPolicyWeight(), PROTECTED, PROBATION);
1 ✔
1252
      remaining--;
1 ✔
1253
    }
1 ✔
1254
  }
1 ✔
1255

1256
  /**
1257
   * Moves the node into another region of the eviction policy, adjusting the weighted size of each
1258
   * region that tracks its own.
1259
   */
1260
  @GuardedBy("evictionLock")
1261
  void transfer(Node<K, V> node, long weight, int from, int to) {
1262
    if (from == WINDOW) {
1 ✔
1263
      setWindowWeightedSize(windowWeightedSize() - weight);
1 ✔
1264
      accessOrderWindowDeque().remove(node);
1 ✔
1265
    } else if (from == PROBATION) {
1 ✔
1266
      accessOrderProbationDeque().remove(node);
1 ✔
1267
    } else {
1268
      setMainProtectedWeightedSize(mainProtectedWeightedSize() - weight);
1 ✔
1269
      accessOrderProtectedDeque().remove(node);
1 ✔
1270
    }
1271

1272
    if (to == WINDOW) {
1 ✔
1273
      setWindowWeightedSize(windowWeightedSize() + weight);
1 ✔
1274
      accessOrderWindowDeque().offerLast(node);
1 ✔
1275
      node.makeWindow();
1 ✔
1276
    } else if (to == PROBATION) {
1 ✔
1277
      accessOrderProbationDeque().offerLast(node);
1 ✔
1278
      node.makeMainProbation();
1 ✔
1279
    } else {
1280
      setMainProtectedWeightedSize(mainProtectedWeightedSize() + weight);
1 ✔
1281
      accessOrderProtectedDeque().offerLast(node);
1 ✔
1282
      node.makeMainProtected();
1 ✔
1283
    }
1284
  }
1 ✔
1285

1286
  /**
1287
   * Performs the post-processing work required after a read.
1288
   *
1289
   * @param node the entry in the page replacement policy
1290
   * @param now the current time, in nanoseconds
1291
   * @param recordHit if the hit count should be incremented
1292
   * @return the refreshed value if immediately loaded, else null
1293
   */
1294
  @Nullable V afterRead(Node<K, V> node, long now, boolean recordHit) {
1295
    if (recordHit) {
1 ✔
1296
      statsCounter().recordHits(1);
1 ✔
1297
    }
1298

1299
    boolean delayable = (readBuffer.offer(node) != Buffer.FULL);
1 ✔
1300
    if (shouldDrainBuffers(delayable)) {
1 ✔
1301
      scheduleDrainBuffers();
1 ✔
1302
    }
1303
    return refreshIfNeeded(node, now);
1 ✔
1304
  }
1305

1306
  /** Enables the read buffer (disabled if fastpath until the frequency sketch is initialized). */
1307
  @GuardedBy("evictionLock")
1308
  void recordReads() {
1309
    if (readBuffer == Buffer.<Node<K, V>>disabled()) {
1 ✔
1310
      // The replacement is published without a fence because a new buffer holds only default state
1311
      readBuffer = new BoundedBuffer<>();
1 ✔
1312
    }
1313
  }
1 ✔
1314

1315
  /**
1316
   * Asynchronously refreshes the entry if eligible.
1317
   *
1318
   * @param node the entry in the cache to refresh
1319
   * @param now the current time, in nanoseconds
1320
   * @return the refreshed value if immediately loaded, else null
1321
   */
1322
  @SuppressWarnings("FutureReturnValueIgnored")
1323
  @Nullable V refreshIfNeeded(Node<K, V> node, long now) {
1324
    if (!refreshAfterWrite()) {
1 ✔
1325
      return null;
1 ✔
1326
    }
1327

1328
    K key;
1329
    V oldValue;
1330
    long writeTime = node.getWriteTime();
1 ✔
1331
    long refreshWriteTime = markRefreshing(writeTime);
1 ✔
1332
    ConcurrentMap<Object, CompletableFuture<?>> refreshes;
1333
    if (((now - writeTime) > refreshAfterWriteNanos())
1 ✔
1334
        && ((key = node.getKey()) != null) && ((oldValue = node.getValue()) != null)
1 ✔
1335
        && !isComputingAsync(oldValue) && !isRefreshing(writeTime)
1 ✔
1336
        && !(refreshes = refreshes()).containsKey(node.getKeyReference())
1 ✔
1337
        && node.isAlive() && node.casWriteTime(writeTime, refreshWriteTime)) {
1 ✔
1338
      long[] startTime = new long[1];
1 ✔
1339
      var keyReference = referenceKey(key);
1 ✔
1340
      @SuppressWarnings({"rawtypes", "unchecked"})
1341
      @Nullable CompletableFuture<? extends @Nullable V>[] refreshFuture = new CompletableFuture[1];
1 ✔
1342
      try {
1343
        refreshes.computeIfAbsent(keyReference, k -> {
1 ✔
1344
          if (!node.isAlive() || (node.getWriteTime() != refreshWriteTime)) {
1 ✔
1345
            return nullRef();
1 ✔
1346
          }
1347
          try {
1348
            startTime[0] = statsTicker().read();
1 ✔
1349
            if (isAsync) {
1 ✔
1350
              @SuppressWarnings("unchecked")
1351
              var future = (CompletableFuture<V>) oldValue;
1 ✔
1352
              if (Async.isReady(future)) {
1 ✔
1353
                requireNonNull(cacheLoader);
1 ✔
1354
                var refresh = cacheLoader.asyncReload(key, future.join(), executor);
1 ✔
1355
                refreshFuture[0] = requireNonNull(refresh, "Null future");
1 ✔
1356
              } else {
1 ✔
1357
                // no-op if the future's completion state was modified (e.g. obtrude methods)
1358
                return nullRef();
1 ✔
1359
              }
1360
            } else {
1 ✔
1361
              requireNonNull(cacheLoader);
1 ✔
1362
              var refresh = cacheLoader.asyncReload(key, oldValue, executor);
1 ✔
1363
              refreshFuture[0] = requireNonNull(refresh, "Null future");
1 ✔
1364
            }
1365
            return castNonNull(refreshFuture[0]);
1 ✔
1366
          } catch (InterruptedException e) {
1 ✔
1367
            Thread.currentThread().interrupt();
1 ✔
1368
            logger.log(Level.WARNING, "Exception thrown when submitting refresh task", e);
1 ✔
1369
            return nullRef();
1 ✔
1370
          } catch (Throwable e) {
1 ✔
1371
            logger.log(Level.WARNING, "Exception thrown when submitting refresh task", e);
1 ✔
1372
            return nullRef();
1 ✔
1373
          }
1374
        });
1375
      } finally {
1376
        node.casWriteTime(refreshWriteTime, writeTime);
1 ✔
1377
      }
1378

1379
      if (refreshFuture[0] == null) {
1 ✔
1380
        return null;
1 ✔
1381
      }
1382

1383
      var refreshed = refreshFuture[0].handle((newValue, error) -> {
1 ✔
1384
        long loadTime = statsTicker().read() - startTime[0];
1 ✔
1385
        if (error != null) {
1 ✔
1386
          if (!(error instanceof CancellationException) && !(error instanceof TimeoutException)) {
1 ✔
1387
            logger.log(Level.WARNING, "Exception thrown during refresh", error);
1 ✔
1388
          }
1389
          refreshes.remove(keyReference, refreshFuture[0]);
1 ✔
1390
          statsCounter().recordLoadFailure(loadTime);
1 ✔
1391
          return null;
1 ✔
1392
        }
1393

1394
        @SuppressWarnings("unchecked")
1395
        V value = (isAsync && (newValue != null)) ? (V) refreshFuture[0] : newValue;
1 ✔
1396
        @Nullable RemovalCause[] cause = new RemovalCause[1];
1 ✔
1397
        var hints = new RemapHints();
1 ✔
1398
        hints.quietly = true;
1 ✔
1399
        V result;
1400
        try {
1401
          result = compute(key, (K k, @Nullable V currentValue) -> {
1 ✔
1402
            // Keep the refresh registered until the write clears it to avoid readers from
1403
            // prematurely scheduling another reload
1404
            boolean owned = (refreshes.get(keyReference) == refreshFuture[0]);
1 ✔
1405
            if (currentValue == null) {
1 ✔
1406
              // If the entry is absent then drop this refresh, unless a successor owns the
1407
              // registration, and maybe notify the listener
1408
              if (value != null) {
1 ✔
1409
                cause[0] = RemovalCause.EXPLICIT;
1 ✔
1410
              }
1411
              hints.preserveRefresh = !owned;
1 ✔
1412
              return null;
1 ✔
1413
            }
1414
            if ((currentValue == oldValue) && node.isAlive()
1 ✔
1415
                && (writeTimeOf(node) == writeTime) && owned) {
1 ✔
1416
              // If the entry was not modified while in-flight (no ABA) then replace, refreshing the
1417
              // metadata even when the reloaded value is the same instance
1418
              return value;
1 ✔
1419
            }
1420
            // Otherwise the refresh is discarded. The rejection installs nothing, so the entry's
1421
            // timestamps are preserved rather than restarted at the completion time. A
1422
            // same-instance reload is not a value change, so it is not notified. When a successor
1423
            // refresh owns the registration, leave it intact so this by-key discard cannot steal
1424
            // its token.
1425
            boolean sameInstance = (currentValue == value)
1 ✔
1426
                || (isAsync && (newValue == Async.getIfReady((CompletableFuture<?>) currentValue)));
1 ✔
1427
            if ((value != null) && !sameInstance) {
1 ✔
1428
              cause[0] = RemovalCause.REPLACED;
1 ✔
1429
            }
1430
            hints.preserveTimestamps = true;
1 ✔
1431
            hints.preserveRefresh = !owned;
1 ✔
1432
            return currentValue;
1 ✔
1433
          }, expiry(), /* recordLoad= */ false, /* recordLoadFailure= */ true, hints);
1 ✔
1434
        } catch (Throwable t) {
1 ✔
1435
          logger.log(Level.WARNING, "Exception thrown during refresh", t);
1 ✔
1436
          refreshes.remove(keyReference, refreshFuture[0]);
1 ✔
1437
          statsCounter().recordLoadFailure(loadTime);
1 ✔
1438
          return null;
1 ✔
1439
        }
1 ✔
1440

1441
        if (cause[0] != null) {
1 ✔
1442
          notifyRemoval(key, value, cause[0]);
1 ✔
1443
        }
1444
        if (newValue == null) {
1 ✔
1445
          statsCounter().recordLoadFailure(loadTime);
1 ✔
1446
        } else {
1447
          statsCounter().recordLoadSuccess(loadTime);
1 ✔
1448
        }
1449
        return result;
1 ✔
1450
      });
1451
      return Async.getIfReady(refreshed);
1 ✔
1452
    }
1453

1454
    return null;
1 ✔
1455
  }
1456

1457
  /**
1458
   * Returns the expiration time for the entry after being created.
1459
   *
1460
   * @param key the key of the entry that was created
1461
   * @param value the value of the entry that was created
1462
   * @param expiry the calculator for the expiration time
1463
   * @param now the current time, in nanoseconds
1464
   * @return the expiration time
1465
   */
1466
  long expireAfterCreate(K key, V value, @Nullable Expiry<? super K, ? super V> expiry, long now) {
1467
    if (expiresVariable()) {
1 ✔
1468
      requireNonNull(expiry);
1 ✔
1469
      long duration = Math.max(0L, expiry.expireAfterCreate(key, value, now));
1 ✔
1470
      return expiresAt(now, duration);
1 ✔
1471
    }
1472
    return 0L;
1 ✔
1473
  }
1474

1475
  /**
1476
   * Returns the expiration time for the entry after being updated.
1477
   *
1478
   * @param node the entry in the page replacement policy
1479
   * @param key the key of the entry that was updated
1480
   * @param value the value of the entry that was updated
1481
   * @param expiry the calculator for the expiration time
1482
   * @param now the current time, in nanoseconds
1483
   * @return the expiration time
1484
   */
1485
  long expireAfterUpdate(Node<K, V> node, K key, V value,
1486
      @Nullable Expiry<? super K, ? super V> expiry, long now) {
1487
    if (expiresVariable()) {
1 ✔
1488
      requireNonNull(expiry);
1 ✔
1489
      long currentDuration = Math.max(0L, node.getVariableTime() - now);
1 ✔
1490
      long duration = Math.max(0L, expiry.expireAfterUpdate(key, value, now, currentDuration));
1 ✔
1491
      return expiresAt(now, duration);
1 ✔
1492
    }
1493
    return 0L;
1 ✔
1494
  }
1495

1496
  /**
1497
   * Returns the access time for the entry after a read.
1498
   *
1499
   * @param node the entry in the page replacement policy
1500
   * @param key the key of the entry that was read
1501
   * @param value the value of the entry that was read
1502
   * @param expiry the calculator for the expiration time
1503
   * @param now the current time, in nanoseconds
1504
   * @return the expiration time
1505
   */
1506
  long expireAfterRead(Node<K, V> node, K key, V value, Expiry<K, V> expiry, long now) {
1507
    if (expiresVariable()) {
1 ✔
1508
      long variableTime = node.getVariableTime();
1 ✔
1509
      long currentDuration = Math.max(0L, variableTime - now);
1 ✔
1510
      if (isAsync && (currentDuration > MAXIMUM_EXPIRY)) {
1 ✔
1511
        // expireAfterCreate has not yet set the duration after completion
1512
        return variableTime;
1 ✔
1513
      }
1514
      long duration = Math.max(0L, expiry.expireAfterRead(key, value, now, currentDuration));
1 ✔
1515
      return expiresAt(now, duration);
1 ✔
1516
    }
1517
    return 0L;
1 ✔
1518
  }
1519

1520
  /**
1521
   * Attempts to update the access time for the entry after a read.
1522
   *
1523
   * @param node the entry in the page replacement policy
1524
   * @param key the key of the entry that was read
1525
   * @param value the value of the entry that was read
1526
   * @param expiry the calculator for the expiration time
1527
   * @param now the current time, in nanoseconds
1528
   */
1529
  void tryExpireAfterRead(Node<K, V> node, K key, V value, Expiry<K, V> expiry, long now) {
1530
    if (!expiresVariable()) {
1 ✔
1531
      return;
1 ✔
1532
    }
1533

1534
    long variableTime = node.getVariableTime();
1 ✔
1535
    long currentDuration = Math.max(0L, variableTime - now);
1 ✔
1536
    if (isAsync && (currentDuration > MAXIMUM_EXPIRY)) {
1 ✔
1537
      // expireAfterCreate has not yet set the duration after completion
1538
      return;
1 ✔
1539
    }
1540

1541
    long tolerance = EXPIRE_TOLERANCE;
1 ✔
1542
    long duration = Math.max(0L, expiry.expireAfterRead(key, value, now, currentDuration));
1 ✔
1543
    long expirationTime = expiresAt(now, duration);
1 ✔
1544
    if (((duration <= tolerance) || (Math.abs(expirationTime - variableTime) > tolerance))
1 ✔
1545
        && (node.getValue() == value)) {
1 ✔
1546
      node.casVariableTime(variableTime, expirationTime);
1 ✔
1547
    }
1548
  }
1 ✔
1549

1550
  void setAccessTime(Node<K, V> node, long now) {
1551
    if (!expiresAfterAccess()) {
1 ✔
1552
      return;
1 ✔
1553
    }
1554
    long tolerance = EXPIRE_TOLERANCE;
1 ✔
1555
    long accessTime = node.getAccessTime();
1 ✔
1556
    if ((expiresAfterAccessNanos() <= tolerance) || (Math.abs(now - accessTime) > tolerance)) {
1 ✔
1557
      node.setAccessTime(now);
1 ✔
1558
    }
1559
  }
1 ✔
1560

1561
  void setWriteTime(Node<K, V> node, long now) {
1562
    if (expiresAfterWrite() || refreshAfterWrite()) {
1 ✔
1563
      // The write time reserves its low bit to mark that a refresh is in flight, so a recorded time
1564
      // is always even and an age must be measured against an evenly truncated clock
1565
      node.setWriteTime(toWriteTime(now));
1 ✔
1566
    }
1567
  }
1 ✔
1568

1569
  void setVariableTime(Node<K, V> node, long expirationTime) {
1570
    if (expiresVariable()) {
1 ✔
1571
      node.setVariableTime(expirationTime);
1 ✔
1572
    }
1573
  }
1 ✔
1574

1575
  /** Returns if the entry's expiration is still deferred to its value's completion. */
1576
  boolean isExpirationDeferred(Node<K, V> node, long now) {
1577
    return expiresVariable() && ((node.getVariableTime() - now) > MAXIMUM_EXPIRY);
1 ✔
1578
  }
1579

1580
  /** Returns if the entry's write time would exceed the minimum expiration reorder threshold. */
1581
  boolean exceedsWriteTimeTolerance(Node<K, V> node, long varTime, long now) {
1582
    long variableTime = node.getVariableTime();
1 ✔
1583
    long writeTime = node.getWriteTime();
1 ✔
1584
    long tolerance = EXPIRE_TOLERANCE;
1 ✔
1585
    return
1 ✔
1586
        (expiresAfterWrite()
1 ✔
1587
            && ((expiresAfterWriteNanos() <= tolerance) || (Math.abs(now - writeTime) > tolerance)))
1 ✔
1588
        || (refreshAfterWrite()
1 ✔
1589
            && ((refreshAfterWriteNanos() <= tolerance) || (Math.abs(now - writeTime) > tolerance)))
1 ✔
1590
        || (expiresVariable() && (Math.abs(varTime - variableTime) > tolerance));
1 ✔
1591
  }
1592

1593
  /** Returns the expiration time for a duration. */
1594
  long expiresAt(long now, long duration) {
1595
    return isAsync ? (now + duration) : (now + Math.min(duration, MAXIMUM_EXPIRY));
1 ✔
1596
  }
1597

1598
  /** Returns the time to record on the entry. */
1599
  long expirationTimeFor(@Nullable V value, long now) {
1600
    return isComputingAsync(value) ? (now + ASYNC_EXPIRY) : now;
1 ✔
1601
  }
1602

1603
  /** Returns the time truncated to the write time's granularity. */
1604
  static long toWriteTime(long time) {
1605
    return (time & ~1L);
1 ✔
1606
  }
1607

1608
  /** Returns the entry's write time, with the refresh marker removed. */
1609
  static long writeTimeOf(Node<?, ?> node) {
1610
    return toWriteTime(node.getWriteTime());
1 ✔
1611
  }
1612

1613
  /** Returns if a refresh is in flight for the entry that recorded this write time. */
1614
  static boolean isRefreshing(long writeTime) {
1615
    return ((writeTime & 1L) != 0L);
1 ✔
1616
  }
1617

1618
  /** Returns the write time marked as having a refresh in flight. */
1619
  static long markRefreshing(long writeTime) {
1620
    return (writeTime | 1L);
1 ✔
1621
  }
1622

1623
  /**
1624
   * Performs the post-processing work required after a write.
1625
   *
1626
   * @param task the pending operation to be applied
1627
   */
1628
  void afterWrite(Runnable task) {
1629
    if (writeBuffer.offer(task)) {
1 ✔
1630
      scheduleAfterWrite();
1 ✔
1631
      return;
1 ✔
1632
    }
1633

1634
    // In scenarios where the writing threads cannot make progress then they attempt to provide
1635
    // assistance by performing the eviction work directly. This can resolve cases where the
1636
    // maintenance task is scheduled but not running. That might occur due to all of the executor's
1637
    // threads being busy (perhaps writing into this cache), the write rate greatly exceeds the
1638
    // consuming rate, priority inversion, or if the executor silently discarded the maintenance
1639
    // task. Unfortunately this cannot resolve when the eviction is blocked waiting on a long-
1640
    // running computation due to an eviction listener, the victim is being computed on by a writer,
1641
    // or the victim residing in the same hash bin as a computing entry. In those cases a warning is
1642
    // logged to encourage the application to decouple these computations from the map operations.
1643
    lock();
1 ✔
1644
    try {
1645
      maintenance(task);
1 ✔
1646
    } catch (RuntimeException e) {
1 ✔
1647
      logger.log(Level.ERROR, "Exception thrown when performing the maintenance task", e);
1 ✔
1648
    } finally {
1649
      evictionLock.unlock();
1 ✔
1650
    }
1651
    rescheduleCleanUpIfIncomplete();
1 ✔
1652
  }
1 ✔
1653

1654
  /** Acquires the eviction lock. */
1655
  void lock() {
1656
    @Var long remainingNanos = WARN_AFTER_LOCK_WAIT_NANOS;
1 ✔
1657
    long end = System.nanoTime() + remainingNanos;
1 ✔
1658
    @Var boolean interrupted = false;
1 ✔
1659
    try {
1660
      for (;;) {
1661
        try {
1662
          if (evictionLock.tryLock(remainingNanos, TimeUnit.NANOSECONDS)) {
1 ✔
1663
            return;
1 ✔
1664
          }
1665
          logger.log(Level.WARNING, "The cache is experiencing excessive wait times for acquiring "
1 ✔
1666
              + "the eviction lock. This may indicate that a long-running computation has halted "
1667
              + "eviction when trying to remove the victim entry. Consider using AsyncCache to "
1668
              + "decouple the computation from the map operation.", new TimeoutException());
1669
          evictionLock.lock();
1 ✔
1670
          return;
1 ✔
1671
        } catch (InterruptedException e) {
1 ✔
1672
          remainingNanos = end - System.nanoTime();
1 ✔
1673
          interrupted = true;
1 ✔
1674
        }
1 ✔
1675
      }
1676
    } finally {
1677
      if (interrupted) {
1 ✔
1678
        Thread.currentThread().interrupt();
1 ✔
1679
      }
1680
    }
1681
  }
1682

1683
  /**
1684
   * Conditionally schedules the asynchronous maintenance task after a write operation. If the
1685
   * task status was IDLE or REQUIRED then the maintenance task is scheduled immediately. If it
1686
   * is already processing then it is set to transition to REQUIRED upon completion so that a new
1687
   * execution is triggered by the next operation.
1688
   */
1689
  void scheduleAfterWrite() {
1690
    @Var int drainStatus = drainStatusOpaque();
1 ✔
1691
    for (;;) {
1692
      switch (drainStatus) {
1 ✔
1693
        case IDLE:
1694
          if (casDrainStatus(IDLE, REQUIRED)) {
1 ✔
1695
            scheduleDrainBuffers();
1 ✔
1696
            return;
1 ✔
1697
          }
1698
          drainStatus = drainStatusAcquire();
1 ✔
1699
          continue;
1 ✔
1700
        case REQUIRED:
1701
          scheduleDrainBuffers();
1 ✔
1702
          return;
1 ✔
1703
        case PROCESSING_TO_IDLE:
1704
          if (casDrainStatus(PROCESSING_TO_IDLE, PROCESSING_TO_REQUIRED)) {
1 ✔
1705
            return;
1 ✔
1706
          }
1707
          drainStatus = drainStatusAcquire();
1 ✔
1708
          continue;
1 ✔
1709
        case PROCESSING_TO_REQUIRED:
1710
          return;
1 ✔
1711
        default:
1712
          throw new IllegalStateException("Invalid drain status: " + drainStatus);
1 ✔
1713
      }
1714
    }
1715
  }
1716

1717
  /**
1718
   * Attempts to schedule an asynchronous task to apply the pending operations to the page
1719
   * replacement policy. If the executor rejects the task then it is run directly.
1720
   */
1721
  void scheduleDrainBuffers() {
1722
    if (drainStatusOpaque() >= PROCESSING_TO_IDLE) {
1 ✔
1723
      return;
1 ✔
1724
    }
1725
    if (evictionLock.tryLock()) {
1 ✔
1726
      try {
1727
        int drainStatus = drainStatusOpaque();
1 ✔
1728
        if (drainStatus >= PROCESSING_TO_IDLE) {
1 ✔
1729
          return;
1 ✔
1730
        }
1731
        setDrainStatusRelease(PROCESSING_TO_IDLE);
1 ✔
1732
        executor.execute(drainBuffersTask);
1 ✔
1733
      } catch (Throwable t) {
1 ✔
1734
        logger.log(Level.WARNING, "Exception thrown when submitting maintenance task", t);
1 ✔
1735
        maintenance(/* ignored */ null);
1 ✔
1736
      } finally {
1737
        evictionLock.unlock();
1 ✔
1738
      }
1739
    }
1740
  }
1 ✔
1741

1742
  @Override
1743
  public void cleanUp() {
1744
    try {
1745
      performCleanUp(/* ignored */ null);
1 ✔
1746
    } catch (RuntimeException e) {
1 ✔
1747
      logger.log(Level.ERROR, "Exception thrown when performing the maintenance task", e);
1 ✔
1748
    }
1 ✔
1749
  }
1 ✔
1750

1751
  /**
1752
   * Performs the maintenance work, blocking until the lock is acquired.
1753
   *
1754
   * @param task an additional pending task to run, or {@code null} if not present
1755
   */
1756
  void performCleanUp(@Nullable Runnable task) {
1757
    evictionLock.lock();
1 ✔
1758
    try {
1759
      maintenance(task);
1 ✔
1760
    } finally {
1761
      evictionLock.unlock();
1 ✔
1762
    }
1763
    rescheduleCleanUpIfIncomplete();
1 ✔
1764
  }
1 ✔
1765

1766
  /**
1767
   * If there remains pending operations that were not handled by the prior clean up then try to
1768
   * schedule an asynchronous maintenance task. This may occur due to a concurrent write after the
1769
   * maintenance work had started or if the amortized threshold of work per clean up was reached.
1770
   */
1771
  @SuppressWarnings("resource")
1772
  void rescheduleCleanUpIfIncomplete() {
1773
    if (drainStatusOpaque() != REQUIRED) {
1 ✔
1774
      return;
1 ✔
1775
    }
1776

1777
    // An immediate scheduling cannot be performed on a custom executor because it may use a
1778
    // caller-runs policy. This could cause the caller's penalty to exceed the amortized threshold,
1779
    // e.g. repeated concurrent writes could result in a retry loop.
1780
    if (executor == ForkJoinPool.commonPool()) {
1 ✔
1781
      scheduleDrainBuffers();
1 ✔
1782
      return;
1 ✔
1783
    }
1784

1785
    // If a scheduler was configured then the maintenance can be deferred onto the custom executor
1786
    // and run in the near future. Otherwise, it will be handled due to other cache activity.
1787
    var pacer = pacer();
1 ✔
1788
    if ((pacer != null) && !pacer.isScheduled() && evictionLock.tryLock()) {
1 ✔
1789
      try {
1790
        if ((drainStatusOpaque() == REQUIRED) && !pacer.isScheduled()) {
1 ✔
1791
          pacer.schedule(executor, drainBuffersTask, expirationTicker().read(), Pacer.TOLERANCE);
1 ✔
1792
        }
1793
      } finally {
1794
        evictionLock.unlock();
1 ✔
1795
      }
1796
    }
1797
  }
1 ✔
1798

1799
  /**
1800
   * Performs the pending maintenance work and sets the state flags during processing to avoid
1801
   * excess scheduling attempts. The read buffer, write buffer, and reference queues are drained,
1802
   * followed by expiration, and size-based eviction.
1803
   *
1804
   * @param task an additional pending task to run, or {@code null} if not present
1805
   */
1806
  @GuardedBy("evictionLock")
1807
  void maintenance(@Nullable Runnable task) {
1808
    setDrainStatusRelease(PROCESSING_TO_IDLE);
1 ✔
1809

1810
    try {
1811
      try {
1812
        drainReadBuffer();
1 ✔
1813
        drainWriteBuffer();
1 ✔
1814
      } finally {
1815
        if (task != null) {
1 ✔
1816
          task.run();
1 ✔
1817
        }
1818
      }
1819

1820
      drainKeyReferences();
1 ✔
1821
      drainValueReferences();
1 ✔
1822

1823
      long now = expirationTicker().read();
1 ✔
1824
      expireEntries(now);
1 ✔
1825
      evictEntries(now);
1 ✔
1826

1827
      climb();
1 ✔
1828
    } finally {
1829
      if ((drainStatusOpaque() != PROCESSING_TO_IDLE)
1 ✔
1830
          || !casDrainStatus(PROCESSING_TO_IDLE, IDLE)) {
1 ✔
1831
        setDrainStatusOpaque(REQUIRED);
1 ✔
1832
      }
1833
    }
1834
  }
1 ✔
1835

1836
  /** Drains the weak key references queue. */
1837
  @GuardedBy("evictionLock")
1838
  void drainKeyReferences() {
1839
    if (!collectKeys()) {
1 ✔
1840
      return;
1 ✔
1841
    }
1842
    for (int i = 0; i <= REFERENCE_THRESHOLD; i++) {
1 ✔
1843
      var keyRef = keyReferenceQueue().poll();
1 ✔
1844
      if (keyRef == null) {
1 ✔
1845
        return;
1 ✔
1846
      }
1847
      Node<K, V> node = data.get(keyRef);
1 ✔
1848
      if (node != null) {
1 ✔
1849
        evictEntry(node, RemovalCause.COLLECTED, 0L);
1 ✔
1850
      }
1851
    }
1852
    setDrainStatusOpaque(PROCESSING_TO_REQUIRED);
1 ✔
1853
  }
1 ✔
1854

1855
  /** Drains the weak / soft value references queue. */
1856
  @GuardedBy("evictionLock")
1857
  void drainValueReferences() {
1858
    if (!collectValues()) {
1 ✔
1859
      return;
1 ✔
1860
    }
1861
    for (int i = 0; i <= REFERENCE_THRESHOLD; i++) {
1 ✔
1862
      var valueRef = valueReferenceQueue().poll();
1 ✔
1863
      if (valueRef == null) {
1 ✔
1864
        return;
1 ✔
1865
      }
1866
      @SuppressWarnings("unchecked")
1867
      var ref = (InternalReference<V>) valueRef;
1 ✔
1868
      Node<K, V> node = data.get(ref.getKeyReference());
1 ✔
1869
      if ((node != null) && (valueRef == node.getValueReference())) {
1 ✔
1870
        evictEntry(node, RemovalCause.COLLECTED, 0L);
1 ✔
1871
      }
1872
    }
1873
    setDrainStatusOpaque(PROCESSING_TO_REQUIRED);
1 ✔
1874
  }
1 ✔
1875

1876
  /** Drains the read buffer. */
1877
  @GuardedBy("evictionLock")
1878
  void drainReadBuffer() {
1879
    readBuffer.drainTo(accessPolicy);
1 ✔
1880
  }
1 ✔
1881

1882
  /**
1883
   * Updates the node's location in the page replacement policy and credits the access.
1884
   * <p>
1885
   * A quiet access reorders but does not record the usage with the admission filter or the adaptive
1886
   * climber's sample, as an asynchronous cache's load completion finalizes the entry's weight
1887
   * rather than conveying a usage; counting it once per load skews the admission frequencies and
1888
   * misattributes window hits.
1889
   */
1890
  @GuardedBy("evictionLock")
1891
  void onAccess(Node<K, V> node, Access access) {
1892
    if (evicts()) {
1 ✔
1893
      var keyRef = node.getKeyReferenceOrNull();
1 ✔
1894
      if ((keyRef == null) || !node.isAlive()) {
1 ✔
1895
        return;
1 ✔
1896
      }
1897
      if (access.recordsUsage()) {
1 ✔
1898
        frequencySketch().increment(keyRef);
1 ✔
1899
      }
1900
      boolean inWindow = node.inWindow();
1 ✔
1901
      boolean inProbation = !inWindow && node.inMainProbation();
1 ✔
1902
      if (inWindow) {
1 ✔
1903
        reorder(accessOrderWindowDeque(), node);
1 ✔
1904
      } else if (inProbation) {
1 ✔
1905
        reorderProbation(node);
1 ✔
1906
      } else {
1907
        reorder(accessOrderProtectedDeque(), node);
1 ✔
1908
      }
1909
      if (access.recordsHit() && (node.getPolicyWeight() != 0)) {
1 ✔
1910
        // a zero-weight entry (an in-flight async load, a pinned entry) earns hits with no
1911
        // capacity, so it must not steer the sizing; the sketch still observes the key
1912
        climber().recordHit(inWindow, inProbation);
1 ✔
1913
      }
1914
    } else if (expiresAfterAccess()) {
1 ✔
1915
      reorder(accessOrderWindowDeque(), node);
1 ✔
1916
    }
1917
    if (expiresVariable()) {
1 ✔
1918
      timerWheel().reschedule(node);
1 ✔
1919
    }
1920
  }
1 ✔
1921

1922
  /** Promote the node from probation to protected on an access. */
1923
  @GuardedBy("evictionLock")
1924
  void reorderProbation(Node<K, V> node) {
1925
    if (!accessOrderProbationDeque().contains(node)) {
1 ✔
1926
      // Ignore stale accesses for an entry that is no longer present
1927
      return;
1 ✔
1928
    } else if (node.getPolicyWeight() > mainProtectedMaximum()) {
1 ✔
1929
      reorder(accessOrderProbationDeque(), node);
1 ✔
1930
      return;
1 ✔
1931
    }
1932

1933
    // If the protected space exceeds its maximum, the LRU items are demoted to the probation space.
1934
    // This is deferred to the adaption phase at the end of the maintenance cycle.
1935
    transfer(node, node.getPolicyWeight(), PROBATION, PROTECTED);
1 ✔
1936
  }
1 ✔
1937

1938
  /** Updates the node's location in the policy's deque, unless it moved to a different one. */
1939
  static <K, V> void reorder(LinkedDeque<Node<K, V>> deque, Node<K, V> node, int queueType) {
1940
    // The access-order deques share the entry's link fields, so containment cannot distinguish
1941
    // which one holds it. A reentrant cycle that transferred the entry cannot be detected, so the
1942
    // scan confirms ownership rather than splicing the deque that now holds it.
1943
    if (node.getQueueType() == queueType) {
1 !
1944
      reorder(deque, node);
1 ✔
1945
    }
1946
  }
1 ✔
1947

1948
  /** Updates the node's location in the policy's deque, unless it is no longer linked. */
1949
  static <K, V> void reorder(LinkedDeque<Node<K, V>> deque, Node<K, V> node) {
1950
    // An entry may be scheduled for reordering despite having been removed. This can occur when the
1951
    // entry was concurrently read while a writer was removing it, or when a reentrant maintenance
1952
    // cycle unlinked it while an expiration scan held it. If the entry is no longer linked then it
1953
    // does not need to be processed.
1954
    if (deque.contains(node)) {
1 ✔
1955
      deque.moveToBack(node);
1 ✔
1956
    }
1957
  }
1 ✔
1958

1959
  /** Drains the write buffer. */
1960
  @GuardedBy("evictionLock")
1961
  void drainWriteBuffer() {
1962
    for (int i = 0; i <= WRITE_BUFFER_MAX; i++) {
1 ✔
1963
      Runnable task = writeBuffer.relaxedPoll();
1 ✔
1964
      if (task == null) {
1 ✔
1965
        return;
1 ✔
1966
      }
1967
      task.run();
1 ✔
1968
    }
1969
    setDrainStatusOpaque(PROCESSING_TO_REQUIRED);
1 ✔
1970
  }
1 ✔
1971

1972
  /** Removes the node from the policy's queues. */
1973
  @GuardedBy("evictionLock")
1974
  void unlink(Node<K, V> node) {
1975
    if (node.inWindow() && (evicts() || expiresAfterAccess())) {
1 ✔
1976
      accessOrderWindowDeque().remove(node);
1 ✔
1977
    } else if (evicts()) {
1 ✔
1978
      if (node.inMainProbation()) {
1 ✔
1979
        accessOrderProbationDeque().remove(node);
1 ✔
1980
      } else {
1981
        accessOrderProtectedDeque().remove(node);
1 ✔
1982
      }
1983
    }
1984
    if (expiresAfterWrite()) {
1 ✔
1985
      writeOrderDeque().remove(node);
1 ✔
1986
    } else if (expiresVariable()) {
1 ✔
1987
      timerWheel().deschedule(node);
1 ✔
1988
    }
1989
  }
1 ✔
1990

1991
  /**
1992
   * Atomically transitions the node to the <code>dead</code> state and decrements the
1993
   * <code>weightedSize</code>.
1994
   *
1995
   * @param node the entry in the page replacement policy
1996
   */
1997
  @GuardedBy("evictionLock")
1998
  @SuppressWarnings("SynchronizationOnLocalVariableOrMethodParameter")
1999
  void makeDead(Node<K, V> node) {
2000
    synchronized (node) {
1 ✔
2001
      if (node.isDead()) {
1 ✔
2002
        return;
1 ✔
2003
      }
2004
      if (evicts()) {
1 ✔
2005
        // The node's policy weight may be out of sync due to a pending update waiting to be
2006
        // processed. At this point the node's weight is finalized, so the weight can be safely
2007
        // taken from the node's perspective and the sizes will be adjusted correctly.
2008
        if (node.inWindow()) {
1 ✔
2009
          setWindowWeightedSize(windowWeightedSize() - node.getWeight());
1 ✔
2010
        } else if (node.inMainProtected()) {
1 ✔
2011
          setMainProtectedWeightedSize(mainProtectedWeightedSize() - node.getWeight());
1 ✔
2012
        }
2013
        setWeightedSize(weightedSize() - node.getWeight());
1 ✔
2014
      }
2015
      node.die();
1 ✔
2016
    }
1 ✔
2017
  }
1 ✔
2018

2019
  /** Adds the node to the page replacement policy. */
2020
  final class AddTask implements Runnable {
2021
    final Node<K, V> node;
2022
    final int weight;
2023

2024
    AddTask(Node<K, V> node, int weight) {
1 ✔
2025
      this.weight = weight;
1 ✔
2026
      this.node = node;
1 ✔
2027
    }
1 ✔
2028

2029
    @Override
2030
    @GuardedBy("evictionLock")
2031
    public void run() {
2032
      if (evicts()) {
1 ✔
2033
        setWeightedSize(weightedSize() + weight);
1 ✔
2034
        setWindowWeightedSize(windowWeightedSize() + weight);
1 ✔
2035
        node.setPolicyWeight(node.getPolicyWeight() + weight);
1 ✔
2036

2037
        long maximum = maximum();
1 ✔
2038
        if (weightedSize() >= (maximum >>> 1)) {
1 ✔
2039
          if (weightedSize() > MAXIMUM_CAPACITY) {
1 ✔
2040
            evictEntries(expirationTicker().read());
1 ✔
2041
          } else {
2042
            // Lazily initialize when close to the maximum
2043
            long capacity = isWeighted() ? data.mappingCount() : maximum;
1 ✔
2044
            frequencySketch().ensureCapacity(capacity);
1 ✔
2045
            recordReads();
1 ✔
2046
          }
2047
        }
2048

2049
        climber().recordMiss();
1 ✔
2050
      }
2051

2052
      // ignore out-of-order write operations
2053
      boolean isAlive;
2054
      synchronized (node) {
1 ✔
2055
        isAlive = node.isAlive();
1 ✔
2056
      }
1 ✔
2057
      if (isAlive) {
1 ✔
2058
        if (expiresAfterWrite()) {
1 ✔
2059
          writeOrderDeque().offerLast(node);
1 ✔
2060
        }
2061
        if (expiresVariable()) {
1 ✔
2062
          timerWheel().schedule(node);
1 ✔
2063
        }
2064
        if (evicts()) {
1 ✔
2065
          var keyRef = node.getKeyReferenceOrNull();
1 ✔
2066
          if ((keyRef != null) && node.isAlive()) {
1 ✔
2067
            frequencySketch().increment(keyRef);
1 ✔
2068
          }
2069
          if (weight > windowMaximum()) {
1 ✔
2070
            accessOrderWindowDeque().offerFirst(node);
1 ✔
2071
          } else {
2072
            accessOrderWindowDeque().offerLast(node);
1 ✔
2073
          }
2074
          if (weight > maximum()) {
1 ✔
2075
            evictEntry(node, RemovalCause.SIZE, expirationTicker().read());
1 ✔
2076
          }
2077
        } else if (expiresAfterAccess()) {
1 ✔
2078
          accessOrderWindowDeque().offerLast(node);
1 ✔
2079
        }
2080
      }
2081
    }
1 ✔
2082
  }
2083

2084
  /** Removes a node from the page replacement policy. */
2085
  final class RemovalTask implements Runnable {
2086
    final Node<K, V> node;
2087

2088
    RemovalTask(Node<K, V> node) {
1 ✔
2089
      this.node = node;
1 ✔
2090
    }
1 ✔
2091

2092
    @Override
2093
    @GuardedBy("evictionLock")
2094
    public void run() {
2095
      // add may not have been processed yet
2096
      unlink(node);
1 ✔
2097
      makeDead(node);
1 ✔
2098
    }
1 ✔
2099
  }
2100

2101
  /** Updates the weighted size. */
2102
  final class UpdateTask implements Runnable {
2103
    final int weightDifference;
2104
    final Node<K, V> node;
2105
    final Access access;
2106

2107
    public UpdateTask(Node<K, V> node, int weightDifference) {
2108
      this(node, weightDifference, Access.HIT);
1 ✔
2109
    }
1 ✔
2110

2111
    public UpdateTask(Node<K, V> node, int weightDifference, Access access) {
1 ✔
2112
      this.weightDifference = weightDifference;
1 ✔
2113
      this.access = access;
1 ✔
2114
      this.node = node;
1 ✔
2115
    }
1 ✔
2116

2117
    @Override
2118
    @GuardedBy("evictionLock")
2119
    public void run() {
2120
      if (expiresAfterWrite()) {
1 ✔
2121
        reorder(writeOrderDeque(), node);
1 ✔
2122
      } else if (expiresVariable()) {
1 ✔
2123
        timerWheel().reschedule(node);
1 ✔
2124
      }
2125
      if (evicts()) {
1 ✔
2126
        if (access.recordsMiss()) {
1 ✔
2127
          climber().recordMiss();
1 ✔
2128
        }
2129
        long oldWeightedSize = node.getPolicyWeight();
1 ✔
2130
        node.setPolicyWeight(oldWeightedSize + weightDifference);
1 ✔
2131
        if (node.inWindow()) {
1 ✔
2132
          setWindowWeightedSize(windowWeightedSize() + weightDifference);
1 ✔
2133
          if (node.getPolicyWeight() > maximum()) {
1 ✔
2134
            evictEntry(node, RemovalCause.SIZE, expirationTicker().read());
1 ✔
2135
          } else if (node.getPolicyWeight() <= windowMaximum()) {
1 ✔
2136
            onAccess(node, access);
1 ✔
2137
          } else if (accessOrderWindowDeque().contains(node)) {
1 ✔
2138
            accessOrderWindowDeque().moveToFront(node);
1 ✔
2139
          }
2140
        } else if (node.inMainProbation()) {
1 ✔
2141
            if (node.getPolicyWeight() <= maximum()) {
1 ✔
2142
              onAccess(node, access);
1 ✔
2143
            } else {
2144
              evictEntry(node, RemovalCause.SIZE, expirationTicker().read());
1 ✔
2145
            }
2146
        } else {
2147
          setMainProtectedWeightedSize(mainProtectedWeightedSize() + weightDifference);
1 ✔
2148
          if (node.getPolicyWeight() <= maximum()) {
1 ✔
2149
            onAccess(node, access);
1 ✔
2150
          } else {
2151
            evictEntry(node, RemovalCause.SIZE, expirationTicker().read());
1 ✔
2152
          }
2153
        }
2154

2155
        setWeightedSize(weightedSize() + weightDifference);
1 ✔
2156
        if (weightedSize() > MAXIMUM_CAPACITY) {
1 ✔
2157
          evictEntries(expirationTicker().read());
1 ✔
2158
        }
2159
      } else if (expiresAfterAccess()) {
1 ✔
2160
        onAccess(node, access);
1 ✔
2161
      }
2162
    }
1 ✔
2163
  }
2164

2165
  /** What an access records beyond the reorder. */
2166
  enum Access {
1 ✔
2167
    HIT(/* usage= */ true, /* hit= */ true, /* miss= */ false),
1 ✔
2168
    RELOAD(/* usage= */ true, /* hit= */ false, /* miss= */ true),
1 ✔
2169
    QUIET(/* usage= */ false, /* hit= */ false, /* miss= */ false);
1 ✔
2170

2171
    private final boolean hit;
2172
    private final boolean miss;
2173
    private final boolean usage;
2174

2175
    Access(boolean usage, boolean hit, boolean miss) {
1 ✔
2176
      this.usage = usage;
1 ✔
2177
      this.hit = hit;
1 ✔
2178
      this.miss = miss;
1 ✔
2179
    }
1 ✔
2180

2181
    /** Returns whether the admission filter observes the key. */
2182
    boolean recordsUsage() {
2183
      return usage;
1 ✔
2184
    }
2185

2186
    /** Returns whether the climber's sample earns a hit. */
2187
    boolean recordsHit() {
2188
      return hit;
1 ✔
2189
    }
2190

2191
    /** Returns whether the climber's sample takes the miss an insertion would have recorded. */
2192
    boolean recordsMiss() {
2193
      return miss;
1 ✔
2194
    }
2195
  }
2196

2197
  /* --------------- Concurrent Map Support --------------- */
2198

2199
  @Override
2200
  public boolean isEmpty() {
2201
    return data.isEmpty();
1 ✔
2202
  }
2203

2204
  @Override
2205
  public int size() {
2206
    return data.size();
1 ✔
2207
  }
2208

2209
  @Override
2210
  public long estimatedSize() {
2211
    return data.mappingCount();
1 ✔
2212
  }
2213

2214
  @Override
2215
  public void clear() {
2216
    // Discard all pending refreshes
2217
    var pending = refreshes;
1 ✔
2218
    if (pending != null) {
1 ✔
2219
      pending.clear();
1 ✔
2220
    }
2221

2222
    Deque<Node<K, V>> entries;
2223
    evictionLock.lock();
1 ✔
2224
    try {
2225
      // Discard all pending reads
2226
      readBuffer.drainTo(e -> {});
1 ✔
2227

2228
      // Apply the pending writes
2229
      for (var task = writeBuffer.relaxedPoll(); task != null; task = writeBuffer.relaxedPoll()) {
1 ✔
2230
        task.run();
1 ✔
2231
      }
2232

2233
      // Cancel the scheduled cleanup
2234
      Pacer pacer = pacer();
1 ✔
2235
      if (pacer != null) {
1 ✔
2236
        pacer.cancel();
1 ✔
2237
      }
2238

2239
      // Discard all entries, falling back to one-by-one to avoid excessive lock hold times
2240
      long now = expirationTicker().read();
1 ✔
2241
      int threshold = (WRITE_BUFFER_MAX / 2);
1 ✔
2242
      entries = new ArrayDeque<>(data.values());
1 ✔
2243
      while (writeBuffer.size() < threshold) {
1 ✔
2244
        var node = entries.pollFirst();
1 ✔
2245
        if (node == null) {
1 ✔
2246
          break;
1 ✔
2247
        }
2248
        removeNode(node, now);
1 ✔
2249
      }
1 ✔
2250
    } finally {
2251
      evictionLock.unlock();
1 ✔
2252
    }
2253

2254
    // Remove any stragglers if released early to more aggressively flush incoming writes
2255
    @Var boolean cleanUp = false;
1 ✔
2256
    for (var node : entries) {
1 ✔
2257
      var key = node.getKey();
1 ✔
2258
      if (key == null) {
1 ✔
2259
        cleanUp = true;
1 ✔
2260
      } else {
2261
        remove(key);
1 ✔
2262
      }
2263
    }
1 ✔
2264
    if (collectKeys() && cleanUp) {
1 ✔
2265
      cleanUp();
1 ✔
2266
    } else {
2267
      rescheduleCleanUpIfIncomplete();
1 ✔
2268
    }
2269
  }
1 ✔
2270

2271
  @GuardedBy("evictionLock")
2272
  @SuppressWarnings({"GuardedByChecker", "SynchronizationOnLocalVariableOrMethodParameter"})
2273
  void removeNode(Node<K, V> node, long now) {
2274
    K key = node.getKey();
1 ✔
2275
    var ctx = new EvictContext<V>();
1 ✔
2276
    var keyReference = node.getKeyReference();
1 ✔
2277

2278
    data.computeIfPresent(keyReference, (k, n) -> {
1 ✔
2279
      if (n != node) {
1 ✔
2280
        return n;
1 ✔
2281
      }
2282
      synchronized (node) {
1 ✔
2283
        ctx.value = node.getValue();
1 ✔
2284
        ctx.oldWeight = node.getWeight();
1 ✔
2285

2286
        if ((key == null) || (ctx.value == null)) {
1 ✔
2287
          ctx.cause = RemovalCause.COLLECTED;
1 ✔
2288
        } else if (hasExpired(node, now) && !isComputingAsync(ctx.value)) {
1 !
2289
          ctx.cause = RemovalCause.EXPIRED;
1 ✔
2290
        } else {
2291
          ctx.cause = RemovalCause.EXPLICIT;
1 ✔
2292
        }
2293

2294
        if (ctx.cause.wasEvicted()) {
1 ✔
2295
          notifyEviction(key, ctx.value, ctx.cause);
1 ✔
2296
        }
2297

2298
        discardRefresh(node.getKeyReference());
1 ✔
2299
        node.retire();
1 ✔
2300
        return null;
1 ✔
2301
      }
2302
    });
2303

2304
    unlink(node);
1 ✔
2305

2306
    synchronized (node) {
1 ✔
2307
      logIfAlive(node);
1 ✔
2308
      makeDead(node);
1 ✔
2309
    }
1 ✔
2310

2311
    if (ctx.cause != null) {
1 ✔
2312
      if (ctx.cause.wasEvicted()) {
1 ✔
2313
        statsCounter().recordEviction(ctx.oldWeight, ctx.cause);
1 ✔
2314
      }
2315
      notifyRemoval(key, ctx.value, ctx.cause);
1 ✔
2316
    }
2317
  }
1 ✔
2318

2319
  @Override
2320
  public boolean containsKey(@Nullable Object key) {
2321
    requireNonNull(key);
1 ✔
2322

2323
    Node<K, V> node = data.get(nodeFactory.newLookupKey(key));
1 ✔
2324
    if (node == null) {
1 ✔
2325
      return false;
1 ✔
2326
    }
2327
    boolean expired = hasExpired(node, expirationTicker().read());
1 ✔
2328
    V value = node.getValue();
1 ✔
2329
    if ((value == null) || (expired && !isComputingAsync(value))) {
1 !
2330
      scheduleDrainBuffers();
1 ✔
2331
      return false;
1 ✔
2332
    }
2333
    return true;
1 ✔
2334
  }
2335

2336
  @Override
2337
  public boolean containsValue(@Nullable Object value) {
2338
    requireNonNull(value);
1 ✔
2339

2340
    long now = expirationTicker().read();
1 ✔
2341
    for (Node<K, V> node : data.values()) {
1 ✔
2342
      boolean expired = hasExpired(node, now);
1 ✔
2343
      V nodeValue = node.getValue();
1 ✔
2344
      if ((node.getKey() == null) || (nodeValue == null)
1 ✔
2345
          || (expired && !isComputingAsync(nodeValue))) {
1 !
2346
        scheduleDrainBuffers();
1 ✔
2347
      } else if (node.isAlive() && node.containsValue(value)) {
1 ✔
2348
        return true;
1 ✔
2349
      }
2350
    }
1 ✔
2351
    return false;
1 ✔
2352
  }
2353

2354
  @Override
2355
  public @Nullable V get(@Nullable Object key) {
2356
    requireNonNull(key);
1 ✔
2357
    return getIfPresent(key, /* recordStats= */ false);
1 ✔
2358
  }
2359

2360
  @Override
2361
  public @Nullable V getIfPresent(Object key, boolean recordStats) {
2362
    Node<K, V> node = data.get(nodeFactory.newLookupKey(key));
1 ✔
2363
    if (node == null) {
1 ✔
2364
      if (recordStats) {
1 ✔
2365
        statsCounter().recordMisses(1);
1 ✔
2366
      }
2367
      if (drainStatusOpaque() == REQUIRED) {
1 ✔
2368
        scheduleDrainBuffers();
1 ✔
2369
      }
2370
      return null;
1 ✔
2371
    }
2372

2373
    long now = expirationTicker().read();
1 ✔
2374
    boolean expired = hasExpired(node, now);
1 ✔
2375
    V value = node.getValue();
1 ✔
2376
    if ((value == null) || (expired && !isComputingAsync(value))) {
1 !
2377
      if (recordStats) {
1 ✔
2378
        statsCounter().recordMisses(1);
1 ✔
2379
      }
2380
      scheduleDrainBuffers();
1 ✔
2381
      return null;
1 ✔
2382
    }
2383

2384
    if (expiresAfterRead() && !isComputingAsync(value)) {
1 ✔
2385
      @SuppressWarnings("unchecked")
2386
      var castedKey = (K) key;
1 ✔
2387
      setAccessTime(node, now);
1 ✔
2388
      tryExpireAfterRead(node, castedKey, value, expiry(), now);
1 ✔
2389
    }
2390
    V refreshed = afterRead(node, now, recordStats);
1 ✔
2391
    return (refreshed == null) ? value : refreshed;
1 ✔
2392
  }
2393

2394
  @Override
2395
  public @Nullable V getIfPresentQuietly(Object key) {
2396
    Node<K, V> node = data.get(nodeFactory.newLookupKey(key));
1 ✔
2397
    if (node == null) {
1 ✔
2398
      return null;
1 ✔
2399
    }
2400
    boolean expired = hasExpired(node, expirationTicker().read());
1 ✔
2401
    V value = node.getValue();
1 ✔
2402
    if ((value == null) || (expired && !isComputingAsync(value))) {
1 !
2403
      return null;
1 ✔
2404
    }
2405
    return value;
1 ✔
2406
  }
2407

2408
  /**
2409
   * Returns the key associated with the mapping in this cache, or {@code null} if there is none.
2410
   *
2411
   * @param key the key whose canonical instance is to be returned
2412
   * @return the key used by the mapping, or {@code null} if this cache does not contain a mapping
2413
   *         for the key
2414
   * @throws NullPointerException if the specified key is null
2415
   */
2416
  public @Nullable K getKey(K key) {
2417
    Node<K, V> node = data.get(nodeFactory.newLookupKey(key));
1 ✔
2418
    if (node == null) {
1 ✔
2419
      if (drainStatusOpaque() == REQUIRED) {
1 ✔
2420
        scheduleDrainBuffers();
1 ✔
2421
      }
2422
      return null;
1 ✔
2423
    }
2424
    afterRead(node, /* now= */ 0L, /* recordHit= */ false);
1 ✔
2425
    return node.getKey();
1 ✔
2426
  }
2427

2428
  @Override
2429
  public Map<K, V> getAllPresent(Iterable<? extends K> keys) {
2430
    var result = new LinkedHashMap<K, @Nullable V>(calculateHashMapCapacity(keys));
1 ✔
2431
    for (K key : keys) {
1 ✔
2432
      result.put(key, null);
1 ✔
2433
    }
1 ✔
2434

2435
    @Var boolean drain = false;
1 ✔
2436
    int uniqueKeys = result.size();
1 ✔
2437
    long now = expirationTicker().read();
1 ✔
2438
    for (var iter = result.entrySet().iterator(); iter.hasNext();) {
1 ✔
2439
      var entry = iter.next();
1 ✔
2440
      Node<K, V> node = data.get(nodeFactory.newLookupKey(entry.getKey()));
1 ✔
2441
      if (node == null) {
1 ✔
2442
        iter.remove();
1 ✔
2443
        continue;
1 ✔
2444
      }
2445
      boolean expired = hasExpired(node, now);
1 ✔
2446
      V value = node.getValue();
1 ✔
2447
      if ((value == null) || (expired && !isComputingAsync(value))) {
1 !
2448
        iter.remove();
1 ✔
2449
        drain = true;
1 ✔
2450
      } else {
2451
        setAccessTime(node, now);
1 ✔
2452
        tryExpireAfterRead(node, entry.getKey(), value, expiry(), now);
1 ✔
2453
        V refreshed = afterRead(node, now, /* recordHit= */ false);
1 ✔
2454
        entry.setValue((refreshed == null) ? value : refreshed);
1 ✔
2455
      }
2456
    }
1 ✔
2457
    if (drain) {
1 ✔
2458
      scheduleDrainBuffers();
1 ✔
2459
    }
2460
    statsCounter().recordHits(result.size());
1 ✔
2461
    statsCounter().recordMisses(uniqueKeys - result.size());
1 ✔
2462

2463
    @SuppressWarnings({"NullableProblems", "NullAway"})
2464
    Map<K, V> unmodifiable = Collections.unmodifiableMap(result);
1 ✔
2465
    return unmodifiable;
1 ✔
2466
  }
2467

2468
  @Override
2469
  public void putAll(Map<? extends K, ? extends V> map) {
2470
    map.forEach(this::put);
1 ✔
2471
  }
1 ✔
2472

2473
  @Override
2474
  public @Nullable V put(K key, V value) {
2475
    return put(key, value, expiry(), /* onlyIfAbsent= */ false);
1 ✔
2476
  }
2477

2478
  @Override
2479
  public @Nullable V putIfAbsent(K key, V value) {
2480
    return put(key, value, expiry(), /* onlyIfAbsent= */ true);
1 ✔
2481
  }
2482

2483
  /**
2484
   * Adds a node to the policy and the data store. If an existing node is found, then its value is
2485
   * updated if allowed.
2486
   *
2487
   * @param key key with which the specified value is to be associated
2488
   * @param value value to be associated with the specified key
2489
   * @param expiry the calculator for the write expiration time
2490
   * @param onlyIfAbsent a write is performed only if the key is not already associated with a value
2491
   * @return the prior value in or null if no mapping was found
2492
   */
2493
  @SuppressWarnings("SynchronizationOnLocalVariableOrMethodParameter")
2494
  @Nullable V put(K key, V value, Expiry<K, V> expiry, boolean onlyIfAbsent) {
2495
    requireNonNull(key);
1 ✔
2496
    requireNonNull(value);
1 ✔
2497

2498
    @Var int newWeight = -1;
1 ✔
2499
    @Var Object keyRef = null;
1 ✔
2500
    @Var Node<K, V> node = null;
1 ✔
2501
    Object lookupKey = nodeFactory.newLookupKey(key);
1 ✔
2502
    for (int attempts = 1; ; attempts++) {
1 ✔
2503
      @Var Node<K, V> prior = data.get(lookupKey);
1 ✔
2504
      if (prior == null) {
1 ✔
2505
        if (node == null) {
1 ✔
2506
          if (newWeight < 0) {
1 ✔
2507
            newWeight = weigher.weigh(key, value);
1 ✔
2508
          }
2509
          long now = expirationTicker().read();
1 ✔
2510
          keyRef = nodeFactory.newReferenceKey(key, keyReferenceQueue());
1 ✔
2511
          node = nodeFactory.newNode(keyRef, value, valueReferenceQueue(), newWeight, now);
1 ✔
2512
          long expirationTime = expirationTimeFor(value, now);
1 ✔
2513
          setVariableTime(node, expireAfterCreate(key, value, expiry, now));
1 ✔
2514
          setAccessTime(node, expirationTime);
1 ✔
2515
          setWriteTime(node, expirationTime);
1 ✔
2516
        }
2517
        var newNode = node;
1 ✔
2518
        var newKeyRef = requireNonNull(keyRef);
1 ✔
2519
        prior = (cacheLoader == null)
1 ✔
2520
            ? data.putIfAbsent(newKeyRef, newNode)
1 ✔
2521
            : data.computeIfAbsent(newKeyRef, k -> {
1 ✔
2522
                discardRefresh(k);
1 ✔
2523
                return newNode;
1 ✔
2524
              });
2525
        if ((prior == null) || (prior == node)) {
1 ✔
2526
          afterWrite(new AddTask(node, newWeight));
1 ✔
2527
          return null;
1 ✔
2528
        } else if (onlyIfAbsent) {
1 ✔
2529
          // An optimistic fast path to avoid unnecessary locking
2530
          long now = expirationTicker().read();
1 ✔
2531
          boolean expired = hasExpired(prior, now);
1 ✔
2532
          V currentValue = prior.getValue();
1 ✔
2533
          if ((currentValue != null) && (!expired || isComputingAsync(currentValue))) {
1 !
2534
            if (!isComputingAsync(currentValue)) {
1 ✔
2535
              tryExpireAfterRead(prior, key, currentValue, expiry, now);
1 ✔
2536
              setAccessTime(prior, now);
1 ✔
2537
            }
2538
            afterRead(prior, now, /* recordHit= */ false);
1 ✔
2539
            return currentValue;
1 ✔
2540
          }
2541
        }
2542
      } else if (onlyIfAbsent) {
1 ✔
2543
        // An optimistic fast path to avoid unnecessary locking
2544
        long now = expirationTicker().read();
1 ✔
2545
        boolean expired = hasExpired(prior, now);
1 ✔
2546
        V currentValue = prior.getValue();
1 ✔
2547
        if ((currentValue != null) && (!expired || isComputingAsync(currentValue))) {
1 !
2548
          if (!isComputingAsync(currentValue)) {
1 ✔
2549
            tryExpireAfterRead(prior, key, currentValue, expiry, now);
1 ✔
2550
            setAccessTime(prior, now);
1 ✔
2551
          }
2552
          afterRead(prior, now, /* recordHit= */ false);
1 ✔
2553
          return currentValue;
1 ✔
2554
        }
2555
      }
2556

2557
      // A read may race with the entry's removal, so that after the entry is acquired it may no
2558
      // longer be usable. A retry will reread from the map and either find an absent mapping, a
2559
      // new entry, or a stale entry.
2560
      if (!prior.isAlive()) {
1 ✔
2561
        // A reread of the stale entry may occur if the state transition occurred but the map
2562
        // removal was delayed by a context switch, so that this thread spin waits until resolved.
2563
        if ((attempts & MAX_PUT_SPIN_WAIT_ATTEMPTS) != 0) {
1 ✔
2564
          Thread.onSpinWait();
1 ✔
2565
          continue;
1 ✔
2566
        }
2567

2568
        // If the spin wait attempts are exhausted then fallback to a map computation in order to
2569
        // deschedule this thread until the entry's removal completes. If the key was modified
2570
        // while in the map so that its equals or hashCode changed then the contents may be
2571
        // corrupted, where the cache holds an evicted (dead) entry that could not be removed.
2572
        // That is a violation of the Map contract, so we check that the mapping is in the "alive"
2573
        // state while in the computation.
2574
        data.computeIfPresent(lookupKey, (k, n) -> {
1 ✔
2575
          requireIsAlive(key, n);
1 ✔
2576
          return n;
1 ✔
2577
        });
2578
        continue;
1 ✔
2579
      }
2580

2581
      long now;
2582
      V oldValue;
2583
      long varTime;
2584
      int oldWeight;
2585
      @Var boolean expired = false;
1 ✔
2586
      @Var boolean mayUpdate = true;
1 ✔
2587
      @Var boolean exceedsTolerance = false;
1 ✔
2588
      if (newWeight < 0) {
1 ✔
2589
        newWeight = weigher.weigh(key, value);
1 ✔
2590
      }
2591
      synchronized (prior) {
1 ✔
2592
        if (!prior.isAlive()) {
1 ✔
2593
          continue;
1 ✔
2594
        }
2595
        oldValue = prior.getValue();
1 ✔
2596
        oldWeight = prior.getWeight();
1 ✔
2597
        now = expirationTicker().read();
1 ✔
2598
        if (oldValue == null) {
1 ✔
2599
          varTime = expireAfterCreate(key, value, expiry, now);
1 ✔
2600
          notifyEviction(key, null, RemovalCause.COLLECTED);
1 ✔
2601
        } else if (hasExpired(prior, now) && !isComputingAsync(oldValue)) {
1 !
2602
          expired = true;
1 ✔
2603
          varTime = expireAfterCreate(key, value, expiry, now);
1 ✔
2604
          notifyEviction(key, oldValue, RemovalCause.EXPIRED);
1 ✔
2605
        } else if (onlyIfAbsent) {
1 ✔
2606
          mayUpdate = false;
1 ✔
2607
          varTime = expireAfterRead(prior, key, oldValue, expiry, now);
1 ✔
2608
        } else {
2609
          varTime = expireAfterUpdate(prior, key, value, expiry, now);
1 ✔
2610
        }
2611

2612
        long expirationTime = expirationTimeFor(mayUpdate ? value : oldValue, now);
1 ✔
2613
        if (mayUpdate) {
1 ✔
2614
          exceedsTolerance = exceedsWriteTimeTolerance(prior, varTime, expirationTime);
1 ✔
2615
          prior.setValue(value, valueReferenceQueue());
1 ✔
2616
          prior.setWeight(newWeight);
1 ✔
2617
          if (expired || exceedsTolerance) {
1 ✔
2618
            setWriteTime(prior, expirationTime);
1 ✔
2619
          }
2620

2621
          discardRefresh(prior.getKeyReference());
1 ✔
2622
        }
2623

2624
        setVariableTime(prior, varTime);
1 ✔
2625
        setAccessTime(prior, expirationTime);
1 ✔
2626
      }
1 ✔
2627

2628
      if (expired) {
1 ✔
2629
        statsCounter().recordEviction(oldWeight, RemovalCause.EXPIRED);
1 ✔
2630
        notifyRemoval(key, oldValue, RemovalCause.EXPIRED);
1 ✔
2631
      } else if (oldValue == null) {
1 ✔
2632
        statsCounter().recordEviction(oldWeight, RemovalCause.COLLECTED);
1 ✔
2633
        notifyRemoval(key, /* value= */ null, RemovalCause.COLLECTED);
1 ✔
2634
      } else if (mayUpdate) {
1 ✔
2635
        notifyOnReplace(key, oldValue, value);
1 ✔
2636
      }
2637

2638
      int weightedDifference = mayUpdate ? (newWeight - oldWeight) : 0;
1 ✔
2639
      if ((oldValue == null) || (weightedDifference != 0) || expired) {
1 ✔
2640
        var access = (expired || (oldValue == null)) ? Access.RELOAD : Access.HIT;
1 ✔
2641
        afterWrite(new UpdateTask(prior, weightedDifference, access));
1 ✔
2642
      } else if (!onlyIfAbsent && exceedsTolerance) {
1 ✔
2643
        afterWrite(new UpdateTask(prior, weightedDifference));
1 ✔
2644
      } else {
2645
        afterRead(prior, now, /* recordHit= */ false);
1 ✔
2646
      }
2647

2648
      return expired ? null : oldValue;
1 ✔
2649
    }
2650
  }
2651

2652
  @Override
2653
  @SuppressWarnings("SynchronizationOnLocalVariableOrMethodParameter")
2654
  public @Nullable V remove(@Nullable Object key) {
2655
    requireNonNull(key);
1 ✔
2656

2657
    var ctx = new RemoveContext<K, V>();
1 ✔
2658
    Object lookupKey = nodeFactory.newLookupKey(key);
1 ✔
2659
    data.compute(lookupKey, (k, n) -> {
1 ✔
2660
      if (n == null) {
1 ✔
2661
        discardRefresh(k);
1 ✔
2662
        return null;
1 ✔
2663
      }
2664
      synchronized (n) {
1 ✔
2665
        requireIsAlive(key, n);
1 ✔
2666
        ctx.oldKey = n.getKey();
1 ✔
2667
        ctx.oldValue = n.getValue();
1 ✔
2668
        ctx.oldWeight = n.getWeight();
1 ✔
2669
        RemovalCause actualCause;
2670
        if ((ctx.oldKey == null) || (ctx.oldValue == null)) {
1 ✔
2671
          actualCause = RemovalCause.COLLECTED;
1 ✔
2672
        } else if (hasExpired(n, expirationTicker().read()) && !isComputingAsync(ctx.oldValue)) {
1 !
2673
          actualCause = RemovalCause.EXPIRED;
1 ✔
2674
        } else {
2675
          actualCause = RemovalCause.EXPLICIT;
1 ✔
2676
        }
2677
        if (actualCause.wasEvicted()) {
1 ✔
2678
          notifyEviction(ctx.oldKey, ctx.oldValue, actualCause);
1 ✔
2679
        }
2680
        ctx.cause = actualCause;
1 ✔
2681
        discardRefresh(k);
1 ✔
2682
        ctx.node = n;
1 ✔
2683
        n.retire();
1 ✔
2684
        return null;
1 ✔
2685
      }
2686
    });
2687

2688
    if (ctx.cause != null) {
1 ✔
2689
      afterWrite(new RemovalTask(requireNonNull(ctx.node)));
1 ✔
2690
      if (ctx.cause.wasEvicted()) {
1 ✔
2691
        statsCounter().recordEviction(ctx.oldWeight, ctx.cause);
1 ✔
2692
      }
2693
      notifyRemoval(ctx.oldKey, ctx.oldValue, ctx.cause);
1 ✔
2694
    }
2695
    return (ctx.cause == RemovalCause.EXPLICIT) ? ctx.oldValue : null;
1 ✔
2696
  }
2697

2698
  @Override
2699
  @SuppressWarnings("SynchronizationOnLocalVariableOrMethodParameter")
2700
  public boolean remove(@Nullable Object key, @Nullable Object value) {
2701
    requireNonNull(key);
1 ✔
2702
    if (value == null) {
1 ✔
2703
      return false;
1 ✔
2704
    }
2705

2706
    var ctx = new RemoveContext<K, V>();
1 ✔
2707
    Object lookupKey = nodeFactory.newLookupKey(key);
1 ✔
2708
    data.computeIfPresent(lookupKey, (kR, node) -> {
1 ✔
2709
      synchronized (node) {
1 ✔
2710
        requireIsAlive(key, node);
1 ✔
2711
        ctx.oldKey = node.getKey();
1 ✔
2712
        ctx.oldValue = node.getValue();
1 ✔
2713
        ctx.oldWeight = node.getWeight();
1 ✔
2714
        if ((ctx.oldKey == null) || (ctx.oldValue == null)) {
1 ✔
2715
          ctx.cause = RemovalCause.COLLECTED;
1 ✔
2716
        } else if (hasExpired(node, expirationTicker().read()) && !isComputingAsync(ctx.oldValue)) {
1 !
2717
          ctx.cause = RemovalCause.EXPIRED;
1 ✔
2718
        } else if (node.containsValue(value)) {
1 ✔
2719
          ctx.cause = RemovalCause.EXPLICIT;
1 ✔
2720
        } else {
2721
          return node;
1 ✔
2722
        }
2723
        if (ctx.cause.wasEvicted()) {
1 ✔
2724
          notifyEviction(ctx.oldKey, ctx.oldValue, ctx.cause);
1 ✔
2725
        }
2726
        discardRefresh(kR);
1 ✔
2727
        ctx.node = node;
1 ✔
2728
        node.retire();
1 ✔
2729
        return null;
1 ✔
2730
      }
2731
    });
2732

2733
    if (ctx.node == null) {
1 ✔
2734
      return false;
1 ✔
2735
    }
2736
    var removeCause = requireNonNull(ctx.cause);
1 ✔
2737
    afterWrite(new RemovalTask(ctx.node));
1 ✔
2738
    if (removeCause.wasEvicted()) {
1 ✔
2739
      statsCounter().recordEviction(ctx.oldWeight, removeCause);
1 ✔
2740
    }
2741
    notifyRemoval(ctx.oldKey, ctx.oldValue, removeCause);
1 ✔
2742

2743
    return (removeCause == RemovalCause.EXPLICIT);
1 ✔
2744
  }
2745

2746
  @Override
2747
  @SuppressWarnings("SynchronizationOnLocalVariableOrMethodParameter")
2748
  public @Nullable V replace(K key, V value) {
2749
    requireNonNull(key);
1 ✔
2750
    requireNonNull(value);
1 ✔
2751
    var ctx = new ReplaceContext<K, V>();
1 ✔
2752
    int weight = weigher.weigh(key, value);
1 ✔
2753
    Node<K, V> node = data.computeIfPresent(nodeFactory.newLookupKey(key), (k, n) -> {
1 ✔
2754
      synchronized (n) {
1 ✔
2755
        requireIsAlive(key, n);
1 ✔
2756
        ctx.nodeKey = n.getKey();
1 ✔
2757
        ctx.oldValue = n.getValue();
1 ✔
2758
        ctx.oldWeight = n.getWeight();
1 ✔
2759
        if ((ctx.nodeKey == null) || (ctx.oldValue == null)
1 ✔
2760
            || (hasExpired(n, ctx.now = expirationTicker().read())
1 ✔
2761
                && !isComputingAsync(ctx.oldValue))) {
1 !
2762
          ctx.oldValue = null;
1 ✔
2763
          ctx.garbage = true;
1 ✔
2764
          return n;
1 ✔
2765
        }
2766

2767
        long varTime = expireAfterUpdate(n, key, value, expiry(), ctx.now);
1 ✔
2768
        n.setValue(value, valueReferenceQueue());
1 ✔
2769
        n.setWeight(weight);
1 ✔
2770

2771
        long expirationTime = expirationTimeFor(value, ctx.now);
1 ✔
2772
        ctx.exceedsTolerance = exceedsWriteTimeTolerance(n, varTime, expirationTime);
1 ✔
2773
        if (ctx.exceedsTolerance) {
1 ✔
2774
          setWriteTime(n, expirationTime);
1 ✔
2775
        }
2776
        setAccessTime(n, expirationTime);
1 ✔
2777
        setVariableTime(n, varTime);
1 ✔
2778
        discardRefresh(k);
1 ✔
2779
        return n;
1 ✔
2780
      }
2781
    });
2782

2783
    if ((node == null) || (ctx.nodeKey == null) || (ctx.oldValue == null)) {
1 ✔
2784
      if (ctx.garbage) {
1 ✔
2785
        scheduleDrainBuffers();
1 ✔
2786
      }
2787
      return null;
1 ✔
2788
    }
2789

2790
    int weightedDifference = (weight - ctx.oldWeight);
1 ✔
2791
    if (ctx.exceedsTolerance || (weightedDifference != 0)) {
1 ✔
2792
      afterWrite(new UpdateTask(node, weightedDifference));
1 ✔
2793
    } else {
2794
      afterRead(node, ctx.now, /* recordHit= */ false);
1 ✔
2795
    }
2796

2797
    notifyOnReplace(ctx.nodeKey, ctx.oldValue, value);
1 ✔
2798
    return ctx.oldValue;
1 ✔
2799
  }
2800

2801
  @Override
2802
  @SuppressWarnings("SynchronizationOnLocalVariableOrMethodParameter")
2803
  public boolean replace(K key, V oldValue, V newValue,
2804
      boolean shouldDiscardRefresh, boolean quietly) {
2805
    requireNonNull(key);
1 ✔
2806
    requireNonNull(oldValue);
1 ✔
2807
    requireNonNull(newValue);
1 ✔
2808
    var ctx = new ReplaceContext<K, V>();
1 ✔
2809
    int weight = weigher.weigh(key, newValue);
1 ✔
2810
    Node<K, V> node = data.computeIfPresent(nodeFactory.newLookupKey(key), (k, n) -> {
1 ✔
2811
      synchronized (n) {
1 ✔
2812
        requireIsAlive(key, n);
1 ✔
2813
        ctx.nodeKey = n.getKey();
1 ✔
2814
        ctx.oldValue = n.getValue();
1 ✔
2815
        ctx.oldWeight = n.getWeight();
1 ✔
2816
        if ((ctx.nodeKey == null) || (ctx.oldValue == null)
1 ✔
2817
            || (hasExpired(n, ctx.now = expirationTicker().read())
1 ✔
2818
                && !isComputingAsync(ctx.oldValue))) {
1 !
2819
          ctx.oldValue = null;
1 ✔
2820
          ctx.garbage = true;
1 ✔
2821
          return n;
1 ✔
2822
        } else if (!n.containsValue(oldValue)) {
1 ✔
2823
          ctx.oldValue = null;
1 ✔
2824
          return n;
1 ✔
2825
        }
2826

2827
        // A completion finalizes only what the insertion deferred. Its caller reads the future's
2828
        // readiness before the write, so a value that becomes available in between is evaluated by
2829
        // the insertion, and evaluating it again would charge that creation as an update.
2830
        long varTime = (quietly && !isExpirationDeferred(n, ctx.now))
1 ✔
2831
            ? n.getVariableTime()
1 ✔
2832
            : expireAfterUpdate(n, key, newValue, expiry(), ctx.now);
1 ✔
2833
        n.setValue(newValue, valueReferenceQueue());
1 ✔
2834
        n.setWeight(weight);
1 ✔
2835

2836
        long expirationTime = expirationTimeFor(newValue, ctx.now);
1 ✔
2837
        ctx.exceedsTolerance = exceedsWriteTimeTolerance(n, varTime, expirationTime);
1 ✔
2838
        if (ctx.exceedsTolerance) {
1 ✔
2839
          setWriteTime(n, expirationTime);
1 ✔
2840
        }
2841
        setAccessTime(n, expirationTime);
1 ✔
2842
        setVariableTime(n, varTime);
1 ✔
2843

2844
        if (shouldDiscardRefresh) {
1 ✔
2845
          discardRefresh(k);
1 ✔
2846
        }
2847
      }
1 ✔
2848
      return n;
1 ✔
2849
    });
2850

2851
    if ((node == null) || (ctx.nodeKey == null) || (ctx.oldValue == null)) {
1 ✔
2852
      if (ctx.garbage) {
1 ✔
2853
        scheduleDrainBuffers();
1 ✔
2854
      }
2855
      return false;
1 ✔
2856
    }
2857

2858
    int weightedDifference = (weight - ctx.oldWeight);
1 ✔
2859
    if (ctx.exceedsTolerance || (weightedDifference != 0)) {
1 ✔
2860
      afterWrite(new UpdateTask(node, weightedDifference, quietly ? Access.QUIET : Access.HIT));
1 ✔
2861
    } else if (!quietly) {
1 ✔
2862
      afterRead(node, ctx.now, /* recordHit= */ false);
1 ✔
2863
    }
2864

2865
    notifyOnReplace(ctx.nodeKey, ctx.oldValue, newValue);
1 ✔
2866
    return true;
1 ✔
2867
  }
2868

2869
  @Override
2870
  public void replaceAll(BiFunction<? super K, ? super V, ? extends V> function) {
2871
    requireNonNull(function);
1 ✔
2872

2873
    BiFunction<K, @Nullable V, V> remappingFunction = (key, oldValue) ->
1 ✔
2874
        requireNonNull(function.apply(key, requireNonNull(oldValue)));
1 ✔
2875
    for (K key : keySet()) {
1 ✔
2876
      Object lookupKey = nodeFactory.newLookupKey(key);
1 ✔
2877
      remap(key, lookupKey, remappingFunction, expiry(),
1 ✔
2878
          new ComputeContext<>(expirationTicker().read()), /* computeIfAbsent= */ false);
1 ✔
2879
    }
1 ✔
2880
  }
1 ✔
2881

2882
  @Override
2883
  public @Nullable V computeIfAbsent(K key,
2884
      @Var Function<? super K, ? extends @Nullable V> mappingFunction,
2885
      boolean recordStats, boolean recordLoad) {
2886
    requireNonNull(key);
1 ✔
2887
    requireNonNull(mappingFunction);
1 ✔
2888

2889
    // An optimistic fast path to avoid unnecessary locking
2890
    Node<K, V> node = data.get(nodeFactory.newLookupKey(key));
1 ✔
2891
    long now = expirationTicker().read();
1 ✔
2892
    if (node != null) {
1 ✔
2893
      boolean expired = hasExpired(node, now);
1 ✔
2894
      V value = node.getValue();
1 ✔
2895
      if ((value != null) && (!expired || isComputingAsync(value))) {
1 !
2896
        if (expiresAfterRead() && !isComputingAsync(value)) {
1 ✔
2897
          tryExpireAfterRead(node, key, value, expiry(), now);
1 ✔
2898
          setAccessTime(node, now);
1 ✔
2899
        }
2900
        var refreshed = afterRead(node, now, /* recordHit= */ recordStats);
1 ✔
2901
        return (refreshed == null) ? value : refreshed;
1 ✔
2902
      }
2903
    }
2904
    if (recordStats) {
1 ✔
2905
      mappingFunction = statsAware(mappingFunction, recordLoad);
1 ✔
2906
    }
2907
    Object keyRef = nodeFactory.newReferenceKey(key, keyReferenceQueue());
1 ✔
2908
    return doComputeIfAbsent(key, keyRef, mappingFunction,
1 ✔
2909
        new ComputeContext<>(now), recordStats);
2910
  }
2911

2912
  /** Returns the current value from a computeIfAbsent invocation. */
2913
  @SuppressWarnings("SynchronizationOnLocalVariableOrMethodParameter")
2914
  @Nullable V doComputeIfAbsent(K key, Object keyRef,
2915
      Function<? super K, ? extends @Nullable V> mappingFunction,
2916
      ComputeContext<K, V> ctx, boolean recordStats) {
2917
    Node<K, V> node = data.compute(keyRef, (k, n) -> {
1 ✔
2918
      if (n == null) {
1 ✔
2919
        ctx.newValue = mappingFunction.apply(key);
1 ✔
2920
        if (ctx.newValue == null) {
1 ✔
2921
          discardRefresh(k);
1 ✔
2922
          return null;
1 ✔
2923
        }
2924
        ctx.now = expirationTicker().read();
1 ✔
2925
        ctx.newWeight = weigher.weigh(key, ctx.newValue);
1 ✔
2926
        var created = nodeFactory.newNode(k, ctx.newValue,
1 ✔
2927
            valueReferenceQueue(), ctx.newWeight, ctx.now);
1 ✔
2928
        long expirationTime = expirationTimeFor(ctx.newValue, ctx.now);
1 ✔
2929
        setVariableTime(created, expireAfterCreate(key, ctx.newValue, expiry(), ctx.now));
1 ✔
2930
        setAccessTime(created, expirationTime);
1 ✔
2931
        setWriteTime(created, expirationTime);
1 ✔
2932
        discardRefresh(k);
1 ✔
2933
        return created;
1 ✔
2934
      }
2935

2936
      synchronized (n) {
1 ✔
2937
        requireIsAlive(key, n);
1 ✔
2938
        ctx.nodeKey = n.getKey();
1 ✔
2939
        ctx.oldValue = n.getValue();
1 ✔
2940
        ctx.oldWeight = n.getWeight();
1 ✔
2941
        RemovalCause actualCause;
2942
        if ((ctx.nodeKey == null) || (ctx.oldValue == null)) {
1 ✔
2943
          actualCause = RemovalCause.COLLECTED;
1 ✔
2944
        } else if (hasExpired(n, ctx.now = expirationTicker().read())
1 ✔
2945
            && !isComputingAsync(ctx.oldValue)) {
1 !
2946
          actualCause = RemovalCause.EXPIRED;
1 ✔
2947
        } else {
2948
          return n;
1 ✔
2949
        }
2950

2951
        ctx.cause = actualCause;
1 ✔
2952
        notifyEviction(ctx.nodeKey, ctx.oldValue, actualCause);
1 ✔
2953

2954
        try {
2955
          ctx.newValue = mappingFunction.apply(key);
1 ✔
2956
          if (ctx.newValue == null) {
1 ✔
2957
            discardRefresh(k);
1 ✔
2958
            ctx.removed = n;
1 ✔
2959
            n.retire();
1 ✔
2960
            return null;
1 ✔
2961
          }
2962
          ctx.now = expirationTicker().read();
1 ✔
2963
          ctx.newWeight = weigher.weigh(key, ctx.newValue);
1 ✔
2964
          long varTime = expireAfterCreate(key, ctx.newValue, expiry(), ctx.now);
1 ✔
2965

2966
          n.setValue(ctx.newValue, valueReferenceQueue());
1 ✔
2967
          n.setWeight(ctx.newWeight);
1 ✔
2968

2969
          long expirationTime = expirationTimeFor(ctx.newValue, ctx.now);
1 ✔
2970
          setAccessTime(n, expirationTime);
1 ✔
2971
          setWriteTime(n, expirationTime);
1 ✔
2972
          setVariableTime(n, varTime);
1 ✔
2973
          discardRefresh(k);
1 ✔
2974
          return n;
1 ✔
2975
        } catch (Throwable e) {
1 ✔
2976
          ctx.newValue = null;
1 ✔
2977
          discardRefresh(k);
1 ✔
2978
          ctx.exception = e;
1 ✔
2979
          ctx.removed = n;
1 ✔
2980
          n.retire();
1 ✔
2981
          return null;
1 ✔
2982
        }
2983
      }
2984
    });
2985

2986
    if (ctx.cause != null) {
1 ✔
2987
      statsCounter().recordEviction(ctx.oldWeight, ctx.cause);
1 ✔
2988
      notifyRemoval(ctx.nodeKey, ctx.oldValue, ctx.cause);
1 ✔
2989
    }
2990
    if (node == null) {
1 ✔
2991
      if (ctx.removed != null) {
1 ✔
2992
        afterWrite(new RemovalTask(ctx.removed));
1 ✔
2993
      }
2994
      if (ctx.exception != null) {
1 ✔
2995
        throw toUnchecked(ctx.exception);
1 ✔
2996
      }
2997
      return null;
1 ✔
2998
    }
2999
    if ((ctx.oldValue != null) && (ctx.newValue == null)) {
1 ✔
3000
      if (!isComputingAsync(ctx.oldValue)) {
1 ✔
3001
        tryExpireAfterRead(node, key, ctx.oldValue, expiry(), ctx.now);
1 ✔
3002
        setAccessTime(node, ctx.now);
1 ✔
3003
      }
3004

3005
      var refreshed = afterRead(node, ctx.now, /* recordHit= */ recordStats);
1 ✔
3006
      return (refreshed == null) ? ctx.oldValue : refreshed;
1 ✔
3007
    }
3008
    if ((ctx.oldValue == null) && (ctx.cause == null)) {
1 ✔
3009
      afterWrite(new AddTask(node, ctx.newWeight));
1 ✔
3010
    } else {
3011
      int weightedDifference = (ctx.newWeight - ctx.oldWeight);
1 ✔
3012
      afterWrite(new UpdateTask(node, weightedDifference, Access.RELOAD));
1 ✔
3013
    }
3014

3015
    return ctx.newValue;
1 ✔
3016
  }
3017

3018
  @Override
3019
  public @Nullable V computeIfPresent(K key,
3020
      BiFunction<? super K, ? super V, ? extends @Nullable V> remappingFunction) {
3021
    requireNonNull(key);
1 ✔
3022
    requireNonNull(remappingFunction);
1 ✔
3023

3024
    // An optimistic fast path to avoid unnecessary locking
3025
    Object lookupKey = nodeFactory.newLookupKey(key);
1 ✔
3026
    var node = data.get(lookupKey);
1 ✔
3027
    long now;
3028
    if (node == null) {
1 ✔
3029
      return null;
1 ✔
3030
    }
3031
    boolean expired = hasExpired(node, now = expirationTicker().read());
1 ✔
3032
    V value = node.getValue();
1 ✔
3033
    if ((value == null) || (expired && !isComputingAsync(value))) {
1 !
3034
      scheduleDrainBuffers();
1 ✔
3035
      return null;
1 ✔
3036
    }
3037

3038
    @SuppressWarnings("NullAway")
3039
    BiFunction<? super K, ? super @Nullable V, ? extends @Nullable V> statsAwareRemappingFunction =
1 ✔
3040
        statsAware(remappingFunction, /* recordLoad= */ true, /* recordLoadFailure= */ true);
1 ✔
3041
    return remap(key, lookupKey, statsAwareRemappingFunction,
1 ✔
3042
        expiry(), new ComputeContext<>(now), /* computeIfAbsent= */ false);
1 ✔
3043
  }
3044

3045
  @Override
3046
  public @Nullable V compute(K key,
3047
      BiFunction<? super K, ? super @Nullable V, ? extends @Nullable V> remappingFunction,
3048
      @Nullable Expiry<? super K, ? super V> expiry, boolean recordLoad, boolean recordLoadFailure,
3049
      @Nullable RemapHints hints) {
3050
    requireNonNull(key);
1 ✔
3051
    requireNonNull(remappingFunction);
1 ✔
3052

3053
    Object keyRef = nodeFactory.newReferenceKey(key, keyReferenceQueue());
1 ✔
3054
    var ctx = new ComputeContext<K, V>(expirationTicker().read());
1 ✔
3055
    ctx.hints = hints;
1 ✔
3056
    return remap(key, keyRef, statsAware(remappingFunction, recordLoad, recordLoadFailure),
1 ✔
3057
        expiry, ctx, /* computeIfAbsent= */ true);
3058
  }
3059

3060
  @Override
3061
  public @Nullable V merge(K key, V value,
3062
      BiFunction<? super V, ? super V, ? extends @Nullable V> remappingFunction) {
3063
    requireNonNull(key);
1 ✔
3064
    requireNonNull(value);
1 ✔
3065
    requireNonNull(remappingFunction);
1 ✔
3066

3067
    Object keyRef = nodeFactory.newReferenceKey(key, keyReferenceQueue());
1 ✔
3068
    BiFunction<? super V, ? super V, ? extends @Nullable V> f = statsAware(remappingFunction);
1 ✔
3069
    return remap(key, keyRef,
1 ✔
3070
        (k, oldValue) -> (oldValue == null) ? value : f.apply(oldValue, value), expiry(),
1 ✔
3071
        new ComputeContext<>(expirationTicker().read()), /* computeIfAbsent= */ true);
1 ✔
3072
  }
3073

3074
  /**
3075
   * Attempts to compute a mapping for the specified key and its current mapped value (or
3076
   * {@code null} if there is no current mapping).
3077
   * <p>
3078
   * An entry that has expired or been reference collected is evicted and the computation continues
3079
   * as if the entry had not been present. This method does not pre-screen and does not wrap the
3080
   * remappingFunction to be statistics aware.
3081
   *
3082
   * @param key key with which the specified value is to be associated
3083
   * @param keyRef the key to associate with or a lookup only key if not {@code computeIfAbsent}
3084
   * @param remappingFunction the function to compute a value
3085
   * @param expiry the calculator for the expiration time
3086
   * @param ctx the mutable context for passing state to and from the {@link ConcurrentHashMap}
3087
   *        compute lambda, with {@link ComputeContext#now} set to the current ticker time
3088
   * @param computeIfAbsent if an absent entry can be computed
3089
   * @return the new value associated with the specified key, or null if none
3090
   */
3091
  @SuppressWarnings({"StatementWithEmptyBody", "SynchronizationOnLocalVariableOrMethodParameter"})
3092
  @Nullable V remap(K key, Object keyRef,
3093
      BiFunction<? super K, ? super @Nullable V, ? extends @Nullable V> remappingFunction,
3094
      @Nullable Expiry<? super K, ? super V> expiry,
3095
      ComputeContext<K, V> ctx, boolean computeIfAbsent) {
3096
    Node<K, V> node = data.compute(keyRef, (kr, n) -> {
1 ✔
3097
      if (n == null) {
1 ✔
3098
        if (!computeIfAbsent) {
1 ✔
3099
          return null;
1 ✔
3100
        }
3101
        ctx.newValue = remappingFunction.apply(key, null);
1 ✔
3102
        if (ctx.newValue == null) {
1 ✔
3103
          // A stale refresh whose entry is absent leaves a successor's registration intact
3104
          if ((ctx.hints == null) || !ctx.hints.preserveRefresh) {
1 ✔
3105
            discardRefresh(kr);
1 ✔
3106
          }
3107
          return null;
1 ✔
3108
        }
3109
        try {
3110
          ctx.now = expirationTicker().read();
1 ✔
3111
          ctx.newWeight = weigher.weigh(key, ctx.newValue);
1 ✔
3112
          long varTime = expireAfterCreate(key, ctx.newValue, expiry, ctx.now);
1 ✔
3113
          var created = nodeFactory.newNode(keyRef, ctx.newValue,
1 ✔
3114
              valueReferenceQueue(), ctx.newWeight, ctx.now);
1 ✔
3115

3116
          long expirationTime = expirationTimeFor(ctx.newValue, ctx.now);
1 ✔
3117
          setAccessTime(created, expirationTime);
1 ✔
3118
          setWriteTime(created, expirationTime);
1 ✔
3119
          setVariableTime(created, varTime);
1 ✔
3120
          return created;
1 ✔
3121
        } finally {
3122
          discardRefresh(kr);
1 ✔
3123
        }
3124
      }
3125

3126
      synchronized (n) {
1 ✔
3127
        requireIsAlive(key, n);
1 ✔
3128
        ctx.nodeKey = n.getKey();
1 ✔
3129
        ctx.oldValue = n.getValue();
1 ✔
3130
        ctx.oldWeight = n.getWeight();
1 ✔
3131
        if ((ctx.nodeKey == null) || (ctx.oldValue == null)) {
1 ✔
3132
          ctx.cause = RemovalCause.COLLECTED;
1 ✔
3133
        } else if (hasExpired(n, expirationTicker().read()) && !isComputingAsync(ctx.oldValue)) {
1 !
3134
          ctx.cause = RemovalCause.EXPIRED;
1 ✔
3135
        }
3136
        if (ctx.cause != null) {
1 ✔
3137
          notifyEviction(ctx.nodeKey, ctx.oldValue, ctx.cause);
1 ✔
3138
          if (!computeIfAbsent) {
1 ✔
3139
            discardRefresh(kr);
1 ✔
3140
            ctx.removed = n;
1 ✔
3141
            n.retire();
1 ✔
3142
            return null;
1 ✔
3143
          }
3144
        }
3145

3146
        boolean wasEvicted = (ctx.cause != null);
1 ✔
3147
        try {
3148
          ctx.newValue = remappingFunction.apply(key,
1 ✔
3149
              (ctx.cause == null) ? ctx.oldValue : null);
1 ✔
3150

3151
          if (ctx.newValue == null) {
1 ✔
3152
            if (ctx.cause == null) {
1 ✔
3153
              ctx.cause = RemovalCause.EXPLICIT;
1 ✔
3154
            }
3155
            // A stale refresh whose entry was evicted leaves a successor's registration intact
3156
            if ((ctx.hints == null) || !ctx.hints.preserveRefresh) {
1 ✔
3157
              discardRefresh(kr);
1 ✔
3158
            }
3159
            ctx.removed = n;
1 ✔
3160
            n.retire();
1 ✔
3161
            return null;
1 ✔
3162
          }
3163

3164
          // If the caller flagged a same-instance return as a no-op (e.g., a refresh was rejected
3165
          // and should not touch the entry), skip the metadata updates below.
3166
          if ((ctx.hints != null) && ctx.hints.preserveTimestamps
1 ✔
3167
              && (ctx.newValue == ctx.oldValue) && (ctx.cause == null)) {
3168
            // Skip for query-style callers whose no-op path must leave any in-flight refresh intact
3169
            if (!ctx.hints.preserveRefresh) {
1 ✔
3170
              discardRefresh(kr);
1 ✔
3171
            }
3172
            ctx.unmodified = true;
1 ✔
3173
            return n;
1 ✔
3174
          }
3175

3176
          long varTime;
3177
          ctx.newWeight = weigher.weigh(key, ctx.newValue);
1 ✔
3178
          ctx.now = expirationTicker().read();
1 ✔
3179
          if (ctx.cause == null) {
1 ✔
3180
            if (ctx.newValue != ctx.oldValue) {
1 ✔
3181
              ctx.cause = RemovalCause.REPLACED;
1 ✔
3182
            }
3183
            varTime = expireAfterUpdate(n, key, ctx.newValue, expiry, ctx.now);
1 ✔
3184
          } else {
3185
            varTime = expireAfterCreate(key, ctx.newValue, expiry, ctx.now);
1 ✔
3186
          }
3187

3188
          if (ctx.newValue != ctx.oldValue) {
1 ✔
3189
            n.setValue(ctx.newValue, valueReferenceQueue());
1 ✔
3190
          }
3191
          n.setWeight(ctx.newWeight);
1 ✔
3192

3193
          long expirationTime = expirationTimeFor(ctx.newValue, ctx.now);
1 ✔
3194
          ctx.exceedsTolerance = exceedsWriteTimeTolerance(n, varTime, expirationTime);
1 ✔
3195
          if (((ctx.cause != null) && ctx.cause.wasEvicted()) || ctx.exceedsTolerance) {
1 ✔
3196
            setWriteTime(n, expirationTime);
1 ✔
3197
          }
3198
          setAccessTime(n, expirationTime);
1 ✔
3199
          setVariableTime(n, varTime);
1 ✔
3200
          discardRefresh(kr);
1 ✔
3201
          return n;
1 ✔
3202
        } catch (Throwable e) {
1 ✔
3203
          if ((ctx.hints == null) || !ctx.hints.preserveRefresh) {
1 ✔
3204
            discardRefresh(kr);
1 ✔
3205
          }
3206
          if (!wasEvicted) {
1 ✔
3207
            throw e;
1 ✔
3208
          }
3209
          ctx.newValue = null;
1 ✔
3210
          ctx.exception = e;
1 ✔
3211
          ctx.removed = n;
1 ✔
3212
          n.retire();
1 ✔
3213
          return null;
1 ✔
3214
        }
3215
      }
3216
    });
3217

3218
    if (ctx.cause != null) {
1 ✔
3219
      if (ctx.cause == RemovalCause.REPLACED) {
1 ✔
3220
        requireNonNull(ctx.newValue);
1 ✔
3221
        notifyOnReplace(key, ctx.oldValue, ctx.newValue);
1 ✔
3222
      } else {
3223
        if (ctx.cause.wasEvicted()) {
1 ✔
3224
          statsCounter().recordEviction(ctx.oldWeight, ctx.cause);
1 ✔
3225
        }
3226
        notifyRemoval(ctx.nodeKey, ctx.oldValue, ctx.cause);
1 ✔
3227
      }
3228
    }
3229

3230
    if (ctx.removed != null) {
1 ✔
3231
      afterWrite(new RemovalTask(ctx.removed));
1 ✔
3232
    } else if (node == null) {
1 ✔
3233
      // absent and not computable
3234
    } else if (ctx.unmodified) {
1 ✔
3235
      // caller flagged a same-instance return as a no-op
3236
    } else if ((ctx.oldValue == null) && (ctx.cause == null)) {
1 ✔
3237
      afterWrite(new AddTask(node, ctx.newWeight));
1 ✔
3238
    } else {
3239
      int weightedDifference = ctx.newWeight - ctx.oldWeight;
1 ✔
3240
      boolean quietly = (ctx.hints != null) && ctx.hints.quietly;
1 ✔
3241
      boolean reload = (ctx.cause != null) && ctx.cause.wasEvicted();
1 ✔
3242
      if (ctx.exceedsTolerance || (weightedDifference != 0) || (evicts() && !quietly && reload)) {
1 ✔
3243
        var access = quietly ? Access.QUIET : (reload ? Access.RELOAD : Access.HIT);
1 ✔
3244
        afterWrite(new UpdateTask(node, weightedDifference, access));
1 ✔
3245
      } else {
1 ✔
3246
        if (!quietly) {
1 ✔
3247
          afterRead(node, ctx.now, /* recordHit= */ false);
1 ✔
3248
        }
3249
        if (reload) {
1 ✔
3250
          scheduleDrainBuffers();
1 ✔
3251
        }
3252
      }
3253
    }
3254

3255
    if (ctx.exception != null) {
1 ✔
3256
      throw toUnchecked(ctx.exception);
1 ✔
3257
    }
3258
    return ctx.newValue;
1 ✔
3259
  }
3260

3261
  @Override
3262
  public void forEach(BiConsumer<? super K, ? super V> action) {
3263
    requireNonNull(action);
1 ✔
3264

3265
    for (var iterator = new EntryIterator<>(this); iterator.hasNext();) {
1 ✔
3266
      action.accept(requireNonNull(iterator.key), requireNonNull(iterator.value));
1 ✔
3267
      iterator.advance();
1 ✔
3268
    }
3269
  }
1 ✔
3270

3271
  @Override
3272
  public Set<K> keySet() {
3273
    Set<K> ks = keySet;
1 ✔
3274
    return (ks == null) ? (keySet = new KeySetView<>(this)) : ks;
1 ✔
3275
  }
3276

3277
  @Override
3278
  public Collection<V> values() {
3279
    Collection<V> vs = values;
1 ✔
3280
    return (vs == null) ? (values = new ValuesView<>(this)) : vs;
1 ✔
3281
  }
3282

3283
  @Override
3284
  public Set<Entry<K, V>> entrySet() {
3285
    Set<Entry<K, V>> es = entrySet;
1 ✔
3286
    return (es == null) ? (entrySet = new EntrySetView<>(this)) : es;
1 ✔
3287
  }
3288

3289
  /**
3290
   * Object equality requires reflexive, symmetric, transitive, and consistency properties. Of
3291
   * these, symmetry and consistency require further clarification for how they are upheld.
3292
   * <p>
3293
   * The <i>consistency</i> property between invocations requires that the results are the same if
3294
   * there are no modifications to the information used. Therefore, usages should expect that this
3295
   * operation may return misleading results if either the maps or the data held by them is modified
3296
   * during the execution of this method. This characteristic allows for comparing the map sizes and
3297
   * assuming stable mappings, as done by {@link java.util.AbstractMap}-based maps.
3298
   * <p>
3299
   * The <i>symmetric</i> property requires that the result is the same for all implementations of
3300
   * {@link Map#equals(Object)}. That contract is defined in terms of the stable mappings provided
3301
   * by {@link #entrySet()}, meaning that the {@link #size()} optimization forces that the count is
3302
   * consistent with the mappings when used for an equality check.
3303
   * <p>
3304
   * The cache's {@link #size()} method may include entries that have expired or have been reference
3305
   * collected, but have not yet been removed from the backing map. An iteration over the map may
3306
   * trigger the removal of these dead entries when skipped over during traversal. To ensure
3307
   * consistency and symmetry, usages should call {@link #cleanUp()} before this method while no
3308
   * other concurrent operations are being performed on this cache. This is not done implicitly by
3309
   * {@link #size()} as many usages assume it to be instantaneous and lock-free. As a postcondition
3310
   * the iteration count is verified against the prescreened {@link #size()} so that a concurrent
3311
   * maintenance pass that drops dead entries during traversal is detected and reported as not
3312
   * equal, rather than silently returning {@code true} on the surviving subset.
3313
   */
3314
  @Override
3315
  public boolean equals(@Nullable Object o) {
3316
    if (o == this) {
1 ✔
3317
      return true;
1 ✔
3318
    } else if (!(o instanceof Map)) {
1 ✔
3319
      return false;
1 ✔
3320
    }
3321

3322
    var map = (Map<?, ?>) o;
1 ✔
3323
    int expectedSize = size();
1 ✔
3324
    if (map.size() != expectedSize) {
1 ✔
3325
      return false;
1 ✔
3326
    }
3327

3328
    try {
3329
      @Var int count = 0;
1 ✔
3330
      long now = expirationTicker().read();
1 ✔
3331
      for (var node : data.values()) {
1 ✔
3332
        boolean expired = hasExpired(node, now);
1 ✔
3333
        K key = node.getKey();
1 ✔
3334
        V value = node.getValue();
1 ✔
3335
        if ((key == null) || (value == null) || !node.isAlive()
1 ✔
3336
            || (expired && !isComputingAsync(value))) {
1 !
3337
          scheduleDrainBuffers();
1 ✔
3338
          return false;
1 ✔
3339
        } else {
3340
          var val = map.get(key);
1 ✔
3341
          if ((val == null) || ((val != value) && !val.equals(value))) {
1 ✔
3342
            return false;
1 ✔
3343
          }
3344
        }
3345
        count++;
1 ✔
3346
      }
1 ✔
3347
      return (count == expectedSize);
1 ✔
3348
    } catch (ClassCastException | NullPointerException ignored) {
1 ✔
3349
      return false;
1 ✔
3350
    }
3351
  }
3352

3353
  @Override
3354
  public int hashCode() {
3355
    @Var int hash = 0;
1 ✔
3356
    @Var boolean drain = false;
1 ✔
3357
    long now = expirationTicker().read();
1 ✔
3358
    for (var node : data.values()) {
1 ✔
3359
      boolean expired = hasExpired(node, now);
1 ✔
3360
      K key = node.getKey();
1 ✔
3361
      V value = node.getValue();
1 ✔
3362
      if ((key == null) || (value == null) || !node.isAlive()
1 ✔
3363
          || (expired && !isComputingAsync(value))) {
1 !
3364
        drain = true;
1 ✔
3365
      } else {
3366
        hash += key.hashCode() ^ value.hashCode();
1 ✔
3367
      }
3368
    }
1 ✔
3369
    if (drain) {
1 ✔
3370
      scheduleDrainBuffers();
1 ✔
3371
    }
3372
    return hash;
1 ✔
3373
  }
3374

3375
  @Override
3376
  public String toString() {
3377
    @Var boolean drain = false;
1 ✔
3378
    long now = expirationTicker().read();
1 ✔
3379
    var result = new StringBuilder().append('{');
1 ✔
3380
    for (var node : data.values()) {
1 ✔
3381
      boolean expired = hasExpired(node, now);
1 ✔
3382
      K key = node.getKey();
1 ✔
3383
      V value = node.getValue();
1 ✔
3384
      if ((key == null) || (value == null) || !node.isAlive()
1 ✔
3385
          || (expired && !isComputingAsync(value))) {
1 !
3386
        drain = true;
1 ✔
3387
      } else {
3388
        if (result.length() != 1) {
1 ✔
3389
          result.append(',').append(' ');
1 ✔
3390
        }
3391
        result.append((key == this) ? "(this Map)" : key);
1 ✔
3392
        result.append('=');
1 ✔
3393
        result.append((value == this) ? "(this Map)" : value);
1 ✔
3394
      }
3395
    }
1 ✔
3396
    if (drain) {
1 ✔
3397
      scheduleDrainBuffers();
1 ✔
3398
    }
3399
    return result.append('}').toString();
1 ✔
3400
  }
3401

3402
  /**
3403
   * Returns the computed result from the ordered traversal of the cache entries.
3404
   *
3405
   * @param hottest the coldest or hottest iteration order
3406
   * @param transformer a function that unwraps the value
3407
   * @param mappingFunction the mapping function to compute a value
3408
   * @return the computed value
3409
   */
3410
  @SuppressWarnings("GuardedByChecker")
3411
  <T extends @Nullable Object> T evictionOrder(boolean hottest,
3412
      Function<@Nullable V, @Nullable V> transformer,
3413
      Function<Stream<CacheEntry<K, V>>, T> mappingFunction) {
3414
    Comparator<Node<K, V>> comparator = Comparator.comparingInt(node -> {
1 ✔
3415
      var keyRef = node.getKeyReferenceOrNull();
1 ✔
3416
      return ((keyRef == null) || !node.isAlive()) ? 0 : frequencySketch().frequency(keyRef);
1 ✔
3417
    });
3418
    Iterable<Node<K, V>> iterable;
3419
    if (hottest) {
1 ✔
3420
      iterable = () -> {
1 ✔
3421
        var secondary = PeekingIterator.comparing(
1 ✔
3422
            accessOrderProbationDeque().descendingIterator(),
1 ✔
3423
            accessOrderWindowDeque().descendingIterator(), comparator);
1 ✔
3424
        return PeekingIterator.concat(
1 ✔
3425
            accessOrderProtectedDeque().descendingIterator(), secondary);
1 ✔
3426
      };
3427
    } else {
3428
      iterable = () -> {
1 ✔
3429
        var primary = PeekingIterator.comparing(
1 ✔
3430
            accessOrderWindowDeque().iterator(), accessOrderProbationDeque().iterator(),
1 ✔
3431
            comparator.reversed());
1 ✔
3432
        return PeekingIterator.concat(primary, accessOrderProtectedDeque().iterator());
1 ✔
3433
      };
3434
    }
3435
    return snapshot(iterable, transformer, mappingFunction);
1 ✔
3436
  }
3437

3438
  /**
3439
   * Returns the computed result from the ordered traversal of the cache entries.
3440
   *
3441
   * @param oldest the youngest or oldest iteration order
3442
   * @param transformer a function that unwraps the value
3443
   * @param mappingFunction the mapping function to compute a value
3444
   * @return the computed value
3445
   */
3446
  @SuppressWarnings("GuardedByChecker")
3447
  <T extends @Nullable Object> T expireAfterAccessOrder(boolean oldest,
3448
      Function<@Nullable V, @Nullable V> transformer,
3449
      Function<Stream<CacheEntry<K, V>>, T> mappingFunction) {
3450
    Iterable<Node<K, V>> iterable;
3451
    if (evicts()) {
1 ✔
3452
      iterable = () -> {
1 ✔
3453
        @Var Comparator<Node<K, V>> comparator = (n1, n2) ->
1 ✔
3454
            Long.signum(n1.getAccessTime() - n2.getAccessTime());
1 ✔
3455
        PeekingIterator<Node<K, V>> first;
3456
        PeekingIterator<Node<K, V>> second;
3457
        PeekingIterator<Node<K, V>> third;
3458
        if (oldest) {
1 ✔
3459
          comparator = comparator.reversed();
1 ✔
3460
          first = accessOrderWindowDeque().iterator();
1 ✔
3461
          second = accessOrderProbationDeque().iterator();
1 ✔
3462
          third = accessOrderProtectedDeque().iterator();
1 ✔
3463
        } else {
3464
          first = accessOrderWindowDeque().descendingIterator();
1 ✔
3465
          second = accessOrderProbationDeque().descendingIterator();
1 ✔
3466
          third = accessOrderProtectedDeque().descendingIterator();
1 ✔
3467
        }
3468
        return PeekingIterator.comparing(
1 ✔
3469
            PeekingIterator.comparing(first, second, comparator), third, comparator);
1 ✔
3470
      };
3471
    } else {
3472
      iterable = oldest
1 ✔
3473
          ? accessOrderWindowDeque()
1 ✔
3474
          : accessOrderWindowDeque()::descendingIterator;
1 ✔
3475
    }
3476
    return snapshot(iterable, transformer, mappingFunction);
1 ✔
3477
  }
3478

3479
  /**
3480
   * Returns the computed result from the ordered traversal of the cache entries.
3481
   *
3482
   * @param iterable the supplier of the entries in the cache
3483
   * @param transformer a function that unwraps the value
3484
   * @param mappingFunction the mapping function to compute a value
3485
   * @return the computed value
3486
   */
3487
  @SuppressWarnings("MathClampLong")
3488
  <T extends @Nullable Object> T snapshot(Iterable<Node<K, V>> iterable,
3489
      Function<@Nullable V, @Nullable V> transformer,
3490
      Function<Stream<CacheEntry<K, V>>, T> mappingFunction) {
3491
    requireNonNull(mappingFunction);
1 ✔
3492
    requireNonNull(transformer);
1 ✔
3493
    requireNonNull(iterable);
1 ✔
3494

3495
    evictionLock.lock();
1 ✔
3496
    try {
3497
      maintenance(/* ignored */ null);
1 ✔
3498

3499
      // Obtain the iterator as late as possible for modification count checking
3500
      try (var stream = StreamSupport.stream(Spliterators.spliteratorUnknownSize(
1 ✔
3501
           iterable.iterator(), DISTINCT | ORDERED | NONNULL | IMMUTABLE), /* parallel= */ false)) {
1 ✔
3502
        boolean[] open = { true };
1 ✔
3503
        return mappingFunction.apply(stream.onClose(() -> open[0] = false)
1 ✔
3504
            .peek(ignored -> requireState(open[0], "stream has already been closed"))
1 ✔
3505
            .map(node -> nodeToCacheEntry(node, transformer,
1 ✔
3506
                (int) Math.max(0L, Math.min(node.getPolicyWeight(), Integer.MAX_VALUE))))
1 ✔
3507
            .filter(Objects::nonNull));
1 ✔
3508
      }
3509
    } finally {
3510
      evictionLock.unlock();
1 ✔
3511
      rescheduleCleanUpIfIncomplete();
1 ✔
3512
    }
3513
  }
3514

3515
  /**
3516
   * Returns an entry for the given node if it can be used externally, else null. The weight is
3517
   * the caller's best-effort reading and may lag a concurrent update.
3518
   */
3519
  @Nullable CacheEntry<K, V> nodeToCacheEntry(
3520
      Node<K, V> node, Function<@Nullable V, @Nullable V> transformer, int weight) {
3521
    long now = expirationTicker().read();
1 ✔
3522
    boolean expired = hasExpired(node, now);
1 ✔
3523
    V rawValue = node.getValue();
1 ✔
3524
    if (rawValue == null) {
1 ✔
3525
      return null;
1 ✔
3526
    }
3527
    V value = transformer.apply(rawValue);
1 ✔
3528
    K key = node.getKey();
1 ✔
3529
    if ((key == null) || (value == null) || !node.isAlive()
1 ✔
3530
        || (expired && !isComputingAsync(rawValue))) {
1 !
3531
      return null;
1 ✔
3532
    }
3533

3534
    @Var long expiresAfter = Long.MAX_VALUE;
1 ✔
3535
    if (expiresAfterAccess()) {
1 ✔
3536
      expiresAfter = Math.min(expiresAfter,
1 ✔
3537
          expiresAfterAccessNanos() - (now - node.getAccessTime()));
1 ✔
3538
    }
3539
    if (expiresAfterWrite()) {
1 ✔
3540
      expiresAfter = Math.min(expiresAfter,
1 ✔
3541
          expiresAfterWriteNanos() - (toWriteTime(now) - writeTimeOf(node)));
1 ✔
3542
    }
3543
    if (expiresVariable()) {
1 ✔
3544
      expiresAfter = node.getVariableTime() - now;
1 ✔
3545
    }
3546

3547
    long refreshableAt = refreshAfterWrite()
1 ✔
3548
        ? writeTimeOf(node) + refreshAfterWriteNanos()
1 ✔
3549
        : now + Long.MAX_VALUE;
1 ✔
3550
    return SnapshotEntry.forEntry(key, value, now,
1 ✔
3551
        isWeighted ? weight : 1, now + expiresAfter, refreshableAt);
1 ✔
3552
  }
3553

3554
  /** Mutable context for passing state between a lambda and the caller. */
3555
  static final class EvictContext<V> {
1 ✔
3556
    @Nullable RemovalCause cause;
3557
    @Nullable V value;
3558
    boolean resurrect;
3559
    boolean removed;
3560
    int oldWeight;
3561
  }
3562

3563
  /** Mutable context for passing state between a lambda and the caller. */
3564
  static final class RemoveContext<K, V> {
1 ✔
3565
    @Nullable K oldKey;
3566
    @Nullable V oldValue;
3567
    @Nullable Node<K, V> node;
3568
    @Nullable RemovalCause cause;
3569
    int oldWeight;
3570
  }
3571

3572
  /** Mutable context for passing state between a lambda and the caller. */
3573
  static final class ReplaceContext<K, V> {
1 ✔
3574
    @Nullable K nodeKey;
3575
    @Nullable V oldValue;
3576

3577
    long now;
3578
    int oldWeight;
3579
    boolean garbage;
3580
    boolean exceedsTolerance;
3581
  }
3582

3583
  /** Mutable context for passing state between a lambda and the caller. */
3584
  static final class ComputeContext<K, V> {
3585
    @Nullable K nodeKey;
3586
    @Nullable V oldValue;
3587
    @Nullable V newValue;
3588
    @Nullable Node<K, V> removed;
3589
    @Nullable RemovalCause cause;
3590
    @Nullable Throwable exception;
3591
    @Nullable RemapHints hints;
3592

3593
    long now;
3594
    int oldWeight;
3595
    int newWeight;
3596
    boolean unmodified;
3597
    boolean exceedsTolerance;
3598

3599
    ComputeContext(long now) {
1 ✔
3600
      this.now = now;
1 ✔
3601
    }
1 ✔
3602
  }
3603

3604
  /** A function that produces an unmodifiable map up to the limit in stream order. */
3605
  static final class SizeLimiter<K, V> implements Function<Stream<CacheEntry<K, V>>, Map<K, V>> {
3606
    private final int expectedSize;
3607
    private final long limit;
3608

3609
    SizeLimiter(int expectedSize, long limit) {
1 ✔
3610
      requireArgument(limit >= 0, "limit cannot be negative: %s", limit);
1 ✔
3611
      this.expectedSize = expectedSize;
1 ✔
3612
      this.limit = limit;
1 ✔
3613
    }
1 ✔
3614

3615
    @Override
3616
    public Map<K, V> apply(Stream<CacheEntry<K, V>> stream) {
3617
      var map = new LinkedHashMap<K, V>(calculateHashMapCapacity(expectedSize));
1 ✔
3618
      stream.limit(limit).forEach(entry -> map.put(entry.getKey(), entry.getValue()));
1 ✔
3619
      return Collections.unmodifiableMap(map);
1 ✔
3620
    }
3621
  }
3622

3623
  /** A function that produces an unmodifiable map up to the weighted limit in stream order. */
3624
  static final class WeightLimiter<K, V> implements Function<Stream<CacheEntry<K, V>>, Map<K, V>> {
3625
    private final long weightLimit;
3626

3627
    long weightedSize;
3628

3629
    WeightLimiter(long weightLimit) {
1 ✔
3630
      requireArgument(weightLimit >= 0, "weight limit cannot be negative: %s", weightLimit);
1 ✔
3631
      this.weightLimit = weightLimit;
1 ✔
3632
    }
1 ✔
3633

3634
    @Override
3635
    public Map<K, V> apply(Stream<CacheEntry<K, V>> stream) {
3636
      var map = new LinkedHashMap<K, V>();
1 ✔
3637
      stream.takeWhile(entry -> {
1 ✔
3638
        weightedSize += entry.weight();
1 ✔
3639
        if (weightedSize < 0) {
1 ✔
3640
          weightedSize = Long.MAX_VALUE;
1 ✔
3641
        }
3642
        return (weightedSize <= weightLimit);
1 ✔
3643
      }).forEach(entry -> map.put(entry.getKey(), entry.getValue()));
1 ✔
3644
      return Collections.unmodifiableMap(map);
1 ✔
3645
    }
3646
  }
3647

3648
  /** An adapter to safely externalize the keys. */
3649
  static final class KeySetView<K, V> extends AbstractSet<K> {
3650
    final BoundedLocalCache<K, V> cache;
3651

3652
    KeySetView(BoundedLocalCache<K, V> cache) {
1 ✔
3653
      this.cache = requireNonNull(cache);
1 ✔
3654
    }
1 ✔
3655

3656
    @Override
3657
    public int size() {
3658
      return cache.size();
1 ✔
3659
    }
3660

3661
    @Override
3662
    public void clear() {
3663
      cache.clear();
1 ✔
3664
    }
1 ✔
3665

3666
    @Override
3667
    @SuppressWarnings("SuspiciousMethodCalls")
3668
    public boolean contains(@Nullable Object o) {
3669
      return cache.containsKey(o);
1 ✔
3670
    }
3671

3672
    @Override
3673
    public boolean containsAll(Collection<?> collection) {
3674
      requireNonNull(collection);
1 ✔
3675
      if (collection != this) {
1 ✔
3676
        for (Object o : collection) {
1 ✔
3677
          if ((o == null) || !contains(o)) {
1 ✔
3678
            return false;
1 ✔
3679
          }
3680
        }
1 ✔
3681
      }
3682
      return true;
1 ✔
3683
    }
3684

3685
    @Override
3686
    public boolean removeAll(Collection<?> collection) {
3687
      requireNonNull(collection);
1 ✔
3688
      @Var boolean modified = false;
1 ✔
3689
      if (cache.collectKeys() || ((collection instanceof Set<?>) && (collection.size() > size()))) {
1 ✔
3690
        for (K key : this) {
1 ✔
3691
          if (collection.contains(key)) {
1 ✔
3692
            modified |= remove(key);
1 ✔
3693
          }
3694
        }
1 ✔
3695
      } else {
3696
        for (var item : collection) {
1 ✔
3697
          modified |= (item != null) && remove(item);
1 ✔
3698
        }
1 ✔
3699
      }
3700
      return modified;
1 ✔
3701
    }
3702

3703
    @Override
3704
    public boolean remove(@Nullable Object o) {
3705
      return (cache.remove(o) != null);
1 ✔
3706
    }
3707

3708
    @Override
3709
    public boolean removeIf(Predicate<? super K> filter) {
3710
      requireNonNull(filter);
1 ✔
3711
      @Var boolean modified = false;
1 ✔
3712
      for (K key : this) {
1 ✔
3713
        if (filter.test(key) && remove(key)) {
1 ✔
3714
          modified = true;
1 ✔
3715
        }
3716
      }
1 ✔
3717
      return modified;
1 ✔
3718
    }
3719

3720
    @Override
3721
    public boolean retainAll(Collection<?> collection) {
3722
      requireNonNull(collection);
1 ✔
3723
      @Var boolean modified = false;
1 ✔
3724
      for (K key : this) {
1 ✔
3725
        if (!collection.contains(key) && remove(key)) {
1 ✔
3726
          modified = true;
1 ✔
3727
        }
3728
      }
1 ✔
3729
      return modified;
1 ✔
3730
    }
3731

3732
    @Override
3733
    public Iterator<K> iterator() {
3734
      return new KeyIterator<>(cache);
1 ✔
3735
    }
3736

3737
    @Override
3738
    public Spliterator<K> spliterator() {
3739
      return new KeySpliterator<>(cache);
1 ✔
3740
    }
3741
  }
3742

3743
  /** An adapter to safely externalize the key iterator. */
3744
  static final class KeyIterator<K, V> implements Iterator<K> {
3745
    final EntryIterator<K, V> iterator;
3746

3747
    KeyIterator(BoundedLocalCache<K, V> cache) {
1 ✔
3748
      this.iterator = new EntryIterator<>(cache);
1 ✔
3749
    }
1 ✔
3750

3751
    @Override
3752
    public boolean hasNext() {
3753
      return iterator.hasNext();
1 ✔
3754
    }
3755

3756
    @Override
3757
    public K next() {
3758
      return iterator.nextKey();
1 ✔
3759
    }
3760

3761
    @Override
3762
    public void remove() {
3763
      iterator.remove();
1 ✔
3764
    }
1 ✔
3765
  }
3766

3767
  /** An adapter to safely externalize the key spliterator. */
3768
  static final class KeySpliterator<K, V> implements Spliterator<K> {
3769
    final Spliterator<Node<K, V>> spliterator;
3770
    final BoundedLocalCache<K, V> cache;
3771

3772
    KeySpliterator(BoundedLocalCache<K, V> cache) {
3773
      this(cache, cache.data.values().spliterator());
1 ✔
3774
    }
1 ✔
3775

3776
    KeySpliterator(BoundedLocalCache<K, V> cache, Spliterator<Node<K, V>> spliterator) {
1 ✔
3777
      this.spliterator = requireNonNull(spliterator);
1 ✔
3778
      this.cache = requireNonNull(cache);
1 ✔
3779
    }
1 ✔
3780

3781
    @Override
3782
    public void forEachRemaining(Consumer<? super K> action) {
3783
      requireNonNull(action);
1 ✔
3784
      Consumer<Node<K, V>> consumer = node -> {
1 ✔
3785
        long now = cache.expirationTicker().read();
1 ✔
3786
        boolean expired = cache.hasExpired(node, now);
1 ✔
3787
        K key = node.getKey();
1 ✔
3788
        V value = node.getValue();
1 ✔
3789
        if ((key == null) || (value == null) || (expired && !cache.isComputingAsync(value))) {
1 !
3790
          cache.scheduleDrainBuffers();
1 ✔
3791
        } else if (node.isAlive()) {
1 ✔
3792
          action.accept(key);
1 ✔
3793
        }
3794
      };
1 ✔
3795
      spliterator.forEachRemaining(consumer);
1 ✔
3796
    }
1 ✔
3797

3798
    @Override
3799
    public boolean tryAdvance(Consumer<? super K> action) {
3800
      requireNonNull(action);
1 ✔
3801
      boolean[] advanced = { false };
1 ✔
3802
      Consumer<Node<K, V>> consumer = node -> {
1 ✔
3803
        long now = cache.expirationTicker().read();
1 ✔
3804
        boolean expired = cache.hasExpired(node, now);
1 ✔
3805
        K key = node.getKey();
1 ✔
3806
        V value = node.getValue();
1 ✔
3807
        if ((key == null) || (value == null) || (expired && !cache.isComputingAsync(value))) {
1 !
3808
          cache.scheduleDrainBuffers();
1 ✔
3809
        } else if (node.isAlive()) {
1 ✔
3810
          action.accept(key);
1 ✔
3811
          advanced[0] = true;
1 ✔
3812
        }
3813
      };
1 ✔
3814
      while (spliterator.tryAdvance(consumer)) {
1 ✔
3815
        if (advanced[0]) {
1 ✔
3816
          return true;
1 ✔
3817
        }
3818
      }
3819
      return false;
1 ✔
3820
    }
3821

3822
    @Override
3823
    public @Nullable Spliterator<K> trySplit() {
3824
      Spliterator<Node<K, V>> split = spliterator.trySplit();
1 ✔
3825
      return (split == null) ? null : new KeySpliterator<>(cache, split);
1 ✔
3826
    }
3827

3828
    @Override
3829
    public long estimateSize() {
3830
      return spliterator.estimateSize();
1 ✔
3831
    }
3832

3833
    @Override
3834
    public int characteristics() {
3835
      return DISTINCT | CONCURRENT | NONNULL;
1 ✔
3836
    }
3837
  }
3838

3839
  /** An adapter to safely externalize the values. */
3840
  static final class ValuesView<K, V> extends AbstractCollection<V> {
3841
    final BoundedLocalCache<K, V> cache;
3842

3843
    ValuesView(BoundedLocalCache<K, V> cache) {
1 ✔
3844
      this.cache = requireNonNull(cache);
1 ✔
3845
    }
1 ✔
3846

3847
    @Override
3848
    public int size() {
3849
      return cache.size();
1 ✔
3850
    }
3851

3852
    @Override
3853
    public void clear() {
3854
      cache.clear();
1 ✔
3855
    }
1 ✔
3856

3857
    @Override
3858
    @SuppressWarnings("SuspiciousMethodCalls")
3859
    public boolean contains(@Nullable Object o) {
3860
      return cache.containsValue(o);
1 ✔
3861
    }
3862

3863
    @Override
3864
    public boolean containsAll(Collection<?> collection) {
3865
      requireNonNull(collection);
1 ✔
3866
      if (collection != this) {
1 ✔
3867
        for (Object o : collection) {
1 ✔
3868
          if ((o == null) || !contains(o)) {
1 ✔
3869
            return false;
1 ✔
3870
          }
3871
        }
1 ✔
3872
      }
3873
      return true;
1 ✔
3874
    }
3875

3876
    @Override
3877
    public boolean removeAll(Collection<?> collection) {
3878
      requireNonNull(collection);
1 ✔
3879
      @Var boolean modified = false;
1 ✔
3880
      for (var iterator = new EntryIterator<>(cache); iterator.hasNext();) {
1 ✔
3881
        var key = requireNonNull(iterator.key);
1 ✔
3882
        var value = requireNonNull(iterator.value);
1 ✔
3883
        if (collection.contains(value) && cache.remove(key, value)) {
1 ✔
3884
          modified = true;
1 ✔
3885
        }
3886
        iterator.advance();
1 ✔
3887
      }
1 ✔
3888
      return modified;
1 ✔
3889
    }
3890

3891
    @Override
3892
    public boolean remove(@Nullable Object o) {
3893
      if (o == null) {
1 ✔
3894
        return false;
1 ✔
3895
      }
3896
      for (var iterator = new EntryIterator<>(cache); iterator.hasNext();) {
1 ✔
3897
        var key = requireNonNull(iterator.key);
1 ✔
3898
        var node = requireNonNull(iterator.next);
1 ✔
3899
        var value = requireNonNull(iterator.value);
1 ✔
3900
        if (node.containsValue(o) && cache.remove(key, value)) {
1 ✔
3901
          return true;
1 ✔
3902
        }
3903
        iterator.advance();
1 ✔
3904
      }
1 ✔
3905
      return false;
1 ✔
3906
    }
3907

3908
    @Override
3909
    public boolean removeIf(Predicate<? super V> filter) {
3910
      requireNonNull(filter);
1 ✔
3911
      @Var boolean modified = false;
1 ✔
3912
      for (var iterator = new EntryIterator<>(cache); iterator.hasNext();) {
1 ✔
3913
        var value = requireNonNull(iterator.value);
1 ✔
3914
        if (filter.test(value)) {
1 ✔
3915
          var key = requireNonNull(iterator.key);
1 ✔
3916
          modified |= cache.remove(key, value);
1 ✔
3917
        }
3918
        iterator.advance();
1 ✔
3919
      }
1 ✔
3920
      return modified;
1 ✔
3921
    }
3922

3923
    @Override
3924
    public boolean retainAll(Collection<?> collection) {
3925
      requireNonNull(collection);
1 ✔
3926
      @Var boolean modified = false;
1 ✔
3927
      for (var iterator = new EntryIterator<>(cache); iterator.hasNext();) {
1 ✔
3928
        var key = requireNonNull(iterator.key);
1 ✔
3929
        var value = requireNonNull(iterator.value);
1 ✔
3930
        if (!collection.contains(value) && cache.remove(key, value)) {
1 ✔
3931
          modified = true;
1 ✔
3932
        }
3933
        iterator.advance();
1 ✔
3934
      }
1 ✔
3935
      return modified;
1 ✔
3936
    }
3937

3938
    @Override
3939
    public Iterator<V> iterator() {
3940
      return new ValueIterator<>(cache);
1 ✔
3941
    }
3942

3943
    @Override
3944
    public Spliterator<V> spliterator() {
3945
      return new ValueSpliterator<>(cache);
1 ✔
3946
    }
3947
  }
3948

3949
  /** An adapter to safely externalize the value iterator. */
3950
  static final class ValueIterator<K, V> implements Iterator<V> {
3951
    final EntryIterator<K, V> iterator;
3952

3953
    ValueIterator(BoundedLocalCache<K, V> cache) {
1 ✔
3954
      this.iterator = new EntryIterator<>(cache);
1 ✔
3955
    }
1 ✔
3956

3957
    @Override
3958
    public boolean hasNext() {
3959
      return iterator.hasNext();
1 ✔
3960
    }
3961

3962
    @Override
3963
    public V next() {
3964
      return iterator.nextValue();
1 ✔
3965
    }
3966

3967
    @Override
3968
    public void remove() {
3969
      iterator.remove();
1 ✔
3970
    }
1 ✔
3971
  }
3972

3973
  /** An adapter to safely externalize the value spliterator. */
3974
  static final class ValueSpliterator<K, V> implements Spliterator<V> {
3975
    final Spliterator<Node<K, V>> spliterator;
3976
    final BoundedLocalCache<K, V> cache;
3977

3978
    ValueSpliterator(BoundedLocalCache<K, V> cache) {
3979
      this(cache, cache.data.values().spliterator());
1 ✔
3980
    }
1 ✔
3981

3982
    ValueSpliterator(BoundedLocalCache<K, V> cache, Spliterator<Node<K, V>> spliterator) {
1 ✔
3983
      this.spliterator = requireNonNull(spliterator);
1 ✔
3984
      this.cache = requireNonNull(cache);
1 ✔
3985
    }
1 ✔
3986

3987
    @Override
3988
    public void forEachRemaining(Consumer<? super V> action) {
3989
      requireNonNull(action);
1 ✔
3990
      Consumer<Node<K, V>> consumer = node -> {
1 ✔
3991
        long now = cache.expirationTicker().read();
1 ✔
3992
        boolean expired = cache.hasExpired(node, now);
1 ✔
3993
        K key = node.getKey();
1 ✔
3994
        V value = node.getValue();
1 ✔
3995
        if ((key == null) || (value == null) || (expired && !cache.isComputingAsync(value))) {
1 !
3996
          cache.scheduleDrainBuffers();
1 ✔
3997
        } else if (node.isAlive()) {
1 ✔
3998
          action.accept(value);
1 ✔
3999
        }
4000
      };
1 ✔
4001
      spliterator.forEachRemaining(consumer);
1 ✔
4002
    }
1 ✔
4003

4004
    @Override
4005
    public boolean tryAdvance(Consumer<? super V> action) {
4006
      requireNonNull(action);
1 ✔
4007
      boolean[] advanced = { false };
1 ✔
4008
      Consumer<Node<K, V>> consumer = node -> {
1 ✔
4009
        long now = cache.expirationTicker().read();
1 ✔
4010
        boolean expired = cache.hasExpired(node, now);
1 ✔
4011
        K key = node.getKey();
1 ✔
4012
        V value = node.getValue();
1 ✔
4013
        if ((key == null) || (value == null) || (expired && !cache.isComputingAsync(value))) {
1 !
4014
          cache.scheduleDrainBuffers();
1 ✔
4015
        } else if (node.isAlive()) {
1 ✔
4016
          action.accept(value);
1 ✔
4017
          advanced[0] = true;
1 ✔
4018
        }
4019
      };
1 ✔
4020
      while (spliterator.tryAdvance(consumer)) {
1 ✔
4021
        if (advanced[0]) {
1 ✔
4022
          return true;
1 ✔
4023
        }
4024
      }
4025
      return false;
1 ✔
4026
    }
4027

4028
    @Override
4029
    public @Nullable Spliterator<V> trySplit() {
4030
      Spliterator<Node<K, V>> split = spliterator.trySplit();
1 ✔
4031
      return (split == null) ? null : new ValueSpliterator<>(cache, split);
1 ✔
4032
    }
4033

4034
    @Override
4035
    public long estimateSize() {
4036
      return spliterator.estimateSize();
1 ✔
4037
    }
4038

4039
    @Override
4040
    public int characteristics() {
4041
      return CONCURRENT | NONNULL;
1 ✔
4042
    }
4043
  }
4044

4045
  /** An adapter to safely externalize the entries. */
4046
  static final class EntrySetView<K, V> extends AbstractSet<Entry<K, V>> {
4047
    final BoundedLocalCache<K, V> cache;
4048

4049
    EntrySetView(BoundedLocalCache<K, V> cache) {
1 ✔
4050
      this.cache = requireNonNull(cache);
1 ✔
4051
    }
1 ✔
4052

4053
    @Override
4054
    public int size() {
4055
      return cache.size();
1 ✔
4056
    }
4057

4058
    @Override
4059
    public void clear() {
4060
      cache.clear();
1 ✔
4061
    }
1 ✔
4062

4063
    @Override
4064
    public boolean contains(@Nullable Object o) {
4065
      if (!(o instanceof Entry<?, ?>)) {
1 ✔
4066
        return false;
1 ✔
4067
      }
4068
      var entry = (Entry<?, ?>) o;
1 ✔
4069
      var key = entry.getKey();
1 ✔
4070
      var value = entry.getValue();
1 ✔
4071
      if ((key == null) || (value == null)) {
1 ✔
4072
        return false;
1 ✔
4073
      }
4074
      Node<K, V> node = cache.data.get(cache.nodeFactory.newLookupKey(key));
1 ✔
4075
      if (node == null) {
1 ✔
4076
        return false;
1 ✔
4077
      }
4078
      boolean expired = cache.hasExpired(node, cache.expirationTicker().read());
1 ✔
4079
      V nodeValue = node.getValue();
1 ✔
4080
      return (nodeValue != null) && node.containsValue(value)
1 ✔
4081
          && (!expired || cache.isComputingAsync(nodeValue));
1 !
4082
    }
4083

4084
    @Override
4085
    public boolean removeAll(Collection<?> collection) {
4086
      requireNonNull(collection);
1 ✔
4087
      @Var boolean modified = false;
1 ✔
4088
      if (cache.collectKeys() || cache.collectValues()
1 ✔
4089
          || ((collection instanceof Set<?>) && (collection.size() > size()))) {
1 ✔
4090
        for (var entry : this) {
1 ✔
4091
          if (collection.contains(entry)) {
1 ✔
4092
            modified |= remove(entry);
1 ✔
4093
          }
4094
        }
1 ✔
4095
      } else {
4096
        for (var item : collection) {
1 ✔
4097
          modified |= remove(item);
1 ✔
4098
        }
1 ✔
4099
      }
4100
      return modified;
1 ✔
4101
    }
4102

4103
    @Override
4104
    @SuppressWarnings("SuspiciousMethodCalls")
4105
    public boolean remove(@Nullable Object o) {
4106
      if (!(o instanceof Entry<?, ?>)) {
1 ✔
4107
        return false;
1 ✔
4108
      }
4109
      var entry = (Entry<?, ?>) o;
1 ✔
4110
      var key = entry.getKey();
1 ✔
4111
      return (key != null) && cache.remove(key, entry.getValue());
1 ✔
4112
    }
4113

4114
    @Override
4115
    public boolean removeIf(Predicate<? super Entry<K, V>> filter) {
4116
      requireNonNull(filter);
1 ✔
4117
      @Var boolean modified = false;
1 ✔
4118
      for (var iterator = new EntryIterator<>(cache); iterator.hasNext();) {
1 ✔
4119
        var key = requireNonNull(iterator.key);
1 ✔
4120
        var value = requireNonNull(iterator.value);
1 ✔
4121
        if (filter.test(Map.entry(key, value))) {
1 ✔
4122
          modified |= cache.remove(key, value);
1 ✔
4123
        }
4124
        iterator.advance();
1 ✔
4125
      }
1 ✔
4126
      return modified;
1 ✔
4127
    }
4128

4129
    @Override
4130
    public boolean retainAll(Collection<?> collection) {
4131
      requireNonNull(collection);
1 ✔
4132
      @Var boolean modified = false;
1 ✔
4133
      for (var entry : this) {
1 ✔
4134
        if (!collection.contains(entry) && remove(entry)) {
1 ✔
4135
          modified = true;
1 ✔
4136
        }
4137
      }
1 ✔
4138
      return modified;
1 ✔
4139
    }
4140

4141
    @Override
4142
    public Iterator<Entry<K, V>> iterator() {
4143
      return new EntryIterator<>(cache);
1 ✔
4144
    }
4145

4146
    @Override
4147
    public Spliterator<Entry<K, V>> spliterator() {
4148
      return new EntrySpliterator<>(cache);
1 ✔
4149
    }
4150
  }
4151

4152
  /** An adapter to safely externalize the entry iterator. */
4153
  static final class EntryIterator<K, V> implements Iterator<Entry<K, V>> {
4154
    final BoundedLocalCache<K, V> cache;
4155
    final Iterator<Node<K, V>> iterator;
4156

4157
    @Nullable K key;
4158
    @Nullable V value;
4159
    @Nullable K removalKey;
4160
    @Nullable Node<K, V> next;
4161

4162
    EntryIterator(BoundedLocalCache<K, V> cache) {
1 ✔
4163
      this.iterator = cache.data.values().iterator();
1 ✔
4164
      this.cache = cache;
1 ✔
4165
    }
1 ✔
4166

4167
    @Override
4168
    public boolean hasNext() {
4169
      if (next != null) {
1 ✔
4170
        return true;
1 ✔
4171
      }
4172

4173
      long now = cache.expirationTicker().read();
1 ✔
4174
      while (iterator.hasNext()) {
1 ✔
4175
        next = iterator.next();
1 ✔
4176
        boolean expired = cache.hasExpired(next, now);
1 ✔
4177
        value = next.getValue();
1 ✔
4178
        key = next.getKey();
1 ✔
4179

4180
        boolean evictable = (key == null) || (value == null)
1 ✔
4181
            || (expired && !cache.isComputingAsync(value));
1 !
4182
        if (evictable || !next.isAlive()) {
1 ✔
4183
          if (evictable) {
1 ✔
4184
            cache.scheduleDrainBuffers();
1 ✔
4185
          }
4186
          advance();
1 ✔
4187
          continue;
1 ✔
4188
        }
4189
        return true;
1 ✔
4190
      }
4191
      return false;
1 ✔
4192
    }
4193

4194
    /** Invalidates the current position so that the iterator may compute the next position. */
4195
    void advance() {
4196
      value = null;
1 ✔
4197
      next = null;
1 ✔
4198
      key = null;
1 ✔
4199
    }
1 ✔
4200

4201
    K nextKey() {
4202
      if (!hasNext()) {
1 ✔
4203
        throw new NoSuchElementException();
1 ✔
4204
      }
4205
      removalKey = key;
1 ✔
4206
      advance();
1 ✔
4207
      return requireNonNull(removalKey);
1 ✔
4208
    }
4209

4210
    V nextValue() {
4211
      if (!hasNext()) {
1 ✔
4212
        throw new NoSuchElementException();
1 ✔
4213
      }
4214
      removalKey = key;
1 ✔
4215
      V val = value;
1 ✔
4216
      advance();
1 ✔
4217
      return requireNonNull(val);
1 ✔
4218
    }
4219

4220
    @Override
4221
    public Entry<K, V> next() {
4222
      if (!hasNext()) {
1 ✔
4223
        throw new NoSuchElementException();
1 ✔
4224
      }
4225
      var entry = new WriteThroughEntry<K, @NonNull V>(
1 ✔
4226
          cache, requireNonNull(key), requireNonNull(value));
1 ✔
4227
      removalKey = key;
1 ✔
4228
      advance();
1 ✔
4229
      return entry;
1 ✔
4230
    }
4231

4232
    @Override
4233
    public void remove() {
4234
      if (removalKey == null) {
1 ✔
4235
        throw new IllegalStateException();
1 ✔
4236
      }
4237
      cache.remove(removalKey);
1 ✔
4238
      removalKey = null;
1 ✔
4239
    }
1 ✔
4240
  }
4241

4242
  /** An adapter to safely externalize the entry spliterator. */
4243
  static final class EntrySpliterator<K, V> implements Spliterator<Entry<K, V>> {
4244
    final Spliterator<Node<K, V>> spliterator;
4245
    final BoundedLocalCache<K, V> cache;
4246

4247
    EntrySpliterator(BoundedLocalCache<K, V> cache) {
4248
      this(cache, cache.data.values().spliterator());
1 ✔
4249
    }
1 ✔
4250

4251
    EntrySpliterator(BoundedLocalCache<K, V> cache, Spliterator<Node<K, V>> spliterator) {
1 ✔
4252
      this.spliterator = requireNonNull(spliterator);
1 ✔
4253
      this.cache = requireNonNull(cache);
1 ✔
4254
    }
1 ✔
4255

4256
    @Override
4257
    public void forEachRemaining(Consumer<? super Entry<K, V>> action) {
4258
      requireNonNull(action);
1 ✔
4259
      Consumer<Node<K, V>> consumer = node -> {
1 ✔
4260
        long now = cache.expirationTicker().read();
1 ✔
4261
        boolean expired = cache.hasExpired(node, now);
1 ✔
4262
        K key = node.getKey();
1 ✔
4263
        V value = node.getValue();
1 ✔
4264
        if ((key == null) || (value == null) || (expired && !cache.isComputingAsync(value))) {
1 !
4265
          cache.scheduleDrainBuffers();
1 ✔
4266
        } else if (node.isAlive()) {
1 ✔
4267
          action.accept(new WriteThroughEntry<>(cache, key, value));
1 ✔
4268
        }
4269
      };
1 ✔
4270
      spliterator.forEachRemaining(consumer);
1 ✔
4271
    }
1 ✔
4272

4273
    @Override
4274
    public boolean tryAdvance(Consumer<? super Entry<K, V>> action) {
4275
      requireNonNull(action);
1 ✔
4276
      boolean[] advanced = { false };
1 ✔
4277
      Consumer<Node<K, V>> consumer = node -> {
1 ✔
4278
        long now = cache.expirationTicker().read();
1 ✔
4279
        boolean expired = cache.hasExpired(node, now);
1 ✔
4280
        K key = node.getKey();
1 ✔
4281
        V value = node.getValue();
1 ✔
4282
        if ((key == null) || (value == null) || (expired && !cache.isComputingAsync(value))) {
1 !
4283
          cache.scheduleDrainBuffers();
1 ✔
4284
        } else if (node.isAlive()) {
1 ✔
4285
          action.accept(new WriteThroughEntry<>(cache, key, value));
1 ✔
4286
          advanced[0] = true;
1 ✔
4287
        }
4288
      };
1 ✔
4289
      while (spliterator.tryAdvance(consumer)) {
1 ✔
4290
        if (advanced[0]) {
1 ✔
4291
          return true;
1 ✔
4292
        }
4293
      }
4294
      return false;
1 ✔
4295
    }
4296

4297
    @Override
4298
    public @Nullable Spliterator<Entry<K, V>> trySplit() {
4299
      Spliterator<Node<K, V>> split = spliterator.trySplit();
1 ✔
4300
      return (split == null) ? null : new EntrySpliterator<>(cache, split);
1 ✔
4301
    }
4302

4303
    @Override
4304
    public long estimateSize() {
4305
      return spliterator.estimateSize();
1 ✔
4306
    }
4307

4308
    @Override
4309
    public int characteristics() {
4310
      return DISTINCT | CONCURRENT | NONNULL;
1 ✔
4311
    }
4312
  }
4313

4314
  /** A reusable task that performs the maintenance work; used to avoid wrapping by ForkJoinPool. */
4315
  static final class PerformCleanupTask extends ForkJoinTask<@Nullable Void> implements Runnable {
4316
    private static final long serialVersionUID = 1L;
4317

4318
    final WeakReference<BoundedLocalCache<?, ?>> reference;
4319

4320
    PerformCleanupTask(BoundedLocalCache<?, ?> cache) {
1 ✔
4321
      reference = new WeakReference<>(cache);
1 ✔
4322
    }
1 ✔
4323

4324
    @Override
4325
    protected boolean exec() {
4326
      run();
1 ✔
4327

4328
      // Indicates that the task has not completed to allow subsequent submissions to execute
4329
      return false;
1 ✔
4330
    }
4331

4332
    @Override
4333
    public void run() {
4334
      BoundedLocalCache<?, ?> cache = reference.get();
1 ✔
4335
      if (cache != null) {
1 ✔
4336
        try {
4337
          cache.performCleanUp(/* ignored */ null);
1 ✔
4338
        } catch (Throwable t) {
1 ✔
4339
          logger.log(Level.ERROR, "Exception thrown when performing the maintenance task", t);
1 ✔
4340
        }
1 ✔
4341
      }
4342
    }
1 ✔
4343

4344
    /**
4345
     * This method cannot be ignored due to being final, so a hostile user supplied Executor could
4346
     * forcibly complete the task and halt future executions. There are easier ways to intentionally
4347
     * harm a system, so this is assumed to not happen in practice.
4348
     */
4349
    // public final void quietlyComplete() {}
4350

4351
    @Override public void complete(@Nullable Void value) {}
1 ✔
4352
    @Override public void setRawResult(@Nullable Void value) {}
1 ✔
4353
    @Override public @Nullable Void getRawResult() { return null; }
1 ✔
4354
    @Override public void completeExceptionally(@Nullable Throwable t) {}
1 ✔
4355
    @Override public boolean cancel(boolean mayInterruptIfRunning) { return false; }
1 ✔
4356
  }
4357

4358
  /** Creates a serialization proxy based on the common configuration shared by all cache types. */
4359
  static <K, V> SerializationProxy<K, V> makeSerializationProxy(BoundedLocalCache<?, ?> cache) {
4360
    var proxy = new SerializationProxy<K, V>();
1 ✔
4361
    proxy.weakKeys = cache.collectKeys();
1 ✔
4362
    proxy.weakValues = cache.nodeFactory.weakValues();
1 ✔
4363
    proxy.softValues = cache.nodeFactory.softValues();
1 ✔
4364
    proxy.isRecordingStats = cache.isRecordingStats();
1 ✔
4365
    proxy.evictionListener = cache.evictionListener;
1 ✔
4366
    proxy.removalListener = cache.removalListener();
1 ✔
4367
    proxy.ticker = cache.expirationTicker();
1 ✔
4368
    if (cache.expiresAfterAccess()) {
1 ✔
4369
      proxy.expiresAfterAccessNanos = cache.expiresAfterAccessNanos();
1 ✔
4370
    }
4371
    if (cache.expiresAfterWrite()) {
1 ✔
4372
      proxy.expiresAfterWriteNanos = cache.expiresAfterWriteNanos();
1 ✔
4373
    }
4374
    if (cache.expiresVariable()) {
1 ✔
4375
      proxy.expiry = cache.expiry();
1 ✔
4376
    }
4377
    if (cache.refreshAfterWrite()) {
1 ✔
4378
      proxy.refreshAfterWriteNanos = cache.refreshAfterWriteNanos();
1 ✔
4379
    }
4380
    if (cache.evicts()) {
1 ✔
4381
      if (cache.isWeighted) {
1 ✔
4382
        proxy.weigher = cache.weigher;
1 ✔
4383
        proxy.maximumWeight = cache.maximumAcquire();
1 ✔
4384
      } else {
4385
        proxy.maximumSize = cache.maximumAcquire();
1 ✔
4386
      }
4387
    }
4388
    proxy.cacheLoader = cache.cacheLoader;
1 ✔
4389
    proxy.async = cache.isAsync;
1 ✔
4390
    return proxy;
1 ✔
4391
  }
4392

4393
  /* --------------- Manual Cache --------------- */
4394

4395
  static class BoundedLocalManualCache<K, V> implements LocalManualCache<K, V>, Serializable {
4396
    private static final long serialVersionUID = 1;
4397

4398
    final BoundedLocalCache<K, V> cache;
4399

4400
    @Nullable Policy<K, V> policy;
4401

4402
    BoundedLocalManualCache(Caffeine<K, V> builder) {
4403
      this(builder, null);
1 ✔
4404
    }
1 ✔
4405

4406
    BoundedLocalManualCache(Caffeine<K, V> builder, @Nullable CacheLoader<? super K, V> loader) {
1 ✔
4407
      cache = LocalCacheFactory.newBoundedLocalCache(builder, loader, /* isAsync= */ false);
1 ✔
4408
    }
1 ✔
4409

4410
    @Override
4411
    public final BoundedLocalCache<K, V> cache() {
4412
      return cache;
1 ✔
4413
    }
4414

4415
    @Override
4416
    public final Policy<K, V> policy() {
4417
      if (policy == null) {
1 ✔
4418
        Function<@Nullable V, @Nullable V> identity = v -> v;
1 ✔
4419
        policy = new BoundedPolicy<>(cache, identity, cache.isWeighted);
1 ✔
4420
      }
4421
      return policy;
1 ✔
4422
    }
4423

4424
    private void readObject(ObjectInputStream stream) throws InvalidObjectException {
4425
      throw new InvalidObjectException("Proxy required");
1 ✔
4426
    }
4427

4428
    private Object writeReplace() {
4429
      return makeSerializationProxy(cache);
1 ✔
4430
    }
4431
  }
4432

4433
  @SuppressWarnings({"NullableOptional",
4434
    "OptionalAssignedToNull", "OptionalUsedAsFieldOrParameterType"})
4435
  static final class BoundedPolicy<K, V> implements Policy<K, V> {
4436
    final Function<@Nullable V, @Nullable V> transformer;
4437
    final BoundedLocalCache<K, V> cache;
4438
    final boolean isWeighted;
4439

4440
    @Nullable Optional<Eviction<K, V>> eviction;
4441
    @Nullable Optional<FixedRefresh<K, V>> refreshes;
4442
    @Nullable Optional<FixedExpiration<K, V>> afterWrite;
4443
    @Nullable Optional<FixedExpiration<K, V>> afterAccess;
4444
    @Nullable Optional<VarExpiration<K, V>> variable;
4445

4446
    BoundedPolicy(BoundedLocalCache<K, V> cache,
4447
        Function<@Nullable V, @Nullable V> transformer, boolean isWeighted) {
1 ✔
4448
      this.transformer = transformer;
1 ✔
4449
      this.isWeighted = isWeighted;
1 ✔
4450
      this.cache = cache;
1 ✔
4451
    }
1 ✔
4452

4453
    @Override public boolean isRecordingStats() {
4454
      return cache.isRecordingStats();
1 ✔
4455
    }
4456
    @Override public @Nullable V getIfPresentQuietly(K key) {
4457
      return transformer.apply(cache.getIfPresentQuietly(key));
1 ✔
4458
    }
4459
    @SuppressWarnings("GuardedByChecker")
4460
    @Override public @Nullable CacheEntry<K, V> getEntryIfPresentQuietly(K key) {
4461
      Node<K, V> node = cache.data.get(cache.nodeFactory.newLookupKey(key));
1 ✔
4462
      return (node == null) ? null : cache.nodeToCacheEntry(node, transformer, node.getWeight());
1 ✔
4463
    }
4464
    @SuppressWarnings("Java9CollectionFactory")
4465
    @Override public Map<K, CompletableFuture<V>> refreshes() {
4466
      var refreshes = cache.refreshes;
1 ✔
4467
      if ((refreshes == null) || refreshes.isEmpty()) {
1 ✔
4468
        @SuppressWarnings({"ImmutableMapOf", "RedundantUnmodifiable"})
4469
        Map<K, CompletableFuture<V>> emptyMap = Collections.unmodifiableMap(Collections.emptyMap());
1 ✔
4470
        return emptyMap;
1 ✔
4471
      } else if (cache.collectKeys()) {
1 ✔
4472
        var inFlight = new IdentityHashMap<K, CompletableFuture<V>>(refreshes.size());
1 ✔
4473
        for (var entry : refreshes.entrySet()) {
1 ✔
4474
          @SuppressWarnings("unchecked")
4475
          var key = ((InternalReference<K>) entry.getKey()).get();
1 ✔
4476
          @SuppressWarnings("unchecked")
4477
          var future = (CompletableFuture<V>) entry.getValue();
1 ✔
4478
          if (key != null) {
1 ✔
4479
            inFlight.put(key, future);
1 ✔
4480
          }
4481
        }
1 ✔
4482
        return Collections.unmodifiableMap(inFlight);
1 ✔
4483
      }
4484
      @SuppressWarnings("unchecked")
4485
      var castedRefreshes = (Map<K, CompletableFuture<V>>) (Object) refreshes;
1 ✔
4486
      return Collections.unmodifiableMap(new HashMap<>(castedRefreshes));
1 ✔
4487
    }
4488
    @Override public Optional<Eviction<K, V>> eviction() {
4489
      return cache.evicts()
1 ✔
4490
          ? (eviction == null) ? (eviction = Optional.of(new BoundedEviction())) : eviction
1 ✔
4491
          : Optional.empty();
1 ✔
4492
    }
4493
    @Override public Optional<FixedExpiration<K, V>> expireAfterAccess() {
4494
      if (!cache.expiresAfterAccess()) {
1 ✔
4495
        return Optional.empty();
1 ✔
4496
      }
4497
      return (afterAccess == null)
1 ✔
4498
          ? (afterAccess = Optional.of(new BoundedExpireAfterAccess()))
1 ✔
4499
          : afterAccess;
1 ✔
4500
    }
4501
    @Override public Optional<FixedExpiration<K, V>> expireAfterWrite() {
4502
      if (!cache.expiresAfterWrite()) {
1 ✔
4503
        return Optional.empty();
1 ✔
4504
      }
4505
      return (afterWrite == null)
1 ✔
4506
          ? (afterWrite = Optional.of(new BoundedExpireAfterWrite()))
1 ✔
4507
          : afterWrite;
1 ✔
4508
    }
4509
    @Override public Optional<VarExpiration<K, V>> expireVariably() {
4510
      if (!cache.expiresVariable()) {
1 ✔
4511
        return Optional.empty();
1 ✔
4512
      }
4513
      return (variable == null)
1 ✔
4514
          ? (variable = Optional.of(new BoundedVarExpiration()))
1 ✔
4515
          : variable;
1 ✔
4516
    }
4517
    @Override public Optional<FixedRefresh<K, V>> refreshAfterWrite() {
4518
      if (!cache.refreshAfterWrite()) {
1 ✔
4519
        return Optional.empty();
1 ✔
4520
      }
4521
      return (refreshes == null)
1 ✔
4522
          ? (refreshes = Optional.of(new BoundedRefreshAfterWrite()))
1 ✔
4523
          : refreshes;
1 ✔
4524
    }
4525

4526
    final class BoundedEviction implements Eviction<K, V> {
1 ✔
4527
      @Override public boolean isWeighted() {
4528
        return isWeighted;
1 ✔
4529
      }
4530
      @Override public OptionalInt weightOf(K key) {
4531
        requireNonNull(key);
1 ✔
4532
        if (!isWeighted) {
1 ✔
4533
          return OptionalInt.empty();
1 ✔
4534
        }
4535
        Node<K, V> node = cache.data.get(cache.nodeFactory.newLookupKey(key));
1 ✔
4536
        if (node == null) {
1 ✔
4537
          return OptionalInt.empty();
1 ✔
4538
        }
4539
        boolean expired = cache.hasExpired(node, cache.expirationTicker().read());
1 ✔
4540
        if (expired) {
1 ✔
4541
          return OptionalInt.empty();
1 ✔
4542
        }
4543
        synchronized (node) {
1 ✔
4544
          V value = node.getValue();
1 ✔
4545
          return ((value != null) && node.isAlive() && !cache.isComputingAsync(value))
1 ✔
4546
              ? OptionalInt.of(node.getWeight())
1 ✔
4547
              : OptionalInt.empty();
1 ✔
4548
        }
4549
      }
4550
      @Override public OptionalLong weightedSize() {
4551
        return isWeighted
1 ✔
4552
            ? OptionalLong.of(Math.max(0, cache.weightedSizeAcquire()))
1 ✔
4553
            : OptionalLong.empty();
1 ✔
4554
      }
4555
      @Override public long getMaximum() {
4556
        return cache.maximumAcquire();
1 ✔
4557
      }
4558
      @Override public void setMaximum(long maximum) {
4559
        cache.evictionLock.lock();
1 ✔
4560
        try {
4561
          cache.setMaximumSize(maximum);
1 ✔
4562
          cache.maintenance(/* ignored */ null);
1 ✔
4563
        } finally {
4564
          cache.evictionLock.unlock();
1 ✔
4565
          cache.rescheduleCleanUpIfIncomplete();
1 ✔
4566
        }
4567
      }
1 ✔
4568
      @Override public Map<K, V> coldest(int limit) {
4569
        int expectedSize = Math.min(limit, cache.size());
1 ✔
4570
        var limiter = new SizeLimiter<K, V>(expectedSize, limit);
1 ✔
4571
        return cache.evictionOrder(/* hottest= */ false, transformer, limiter);
1 ✔
4572
      }
4573
      @Override public Map<K, V> coldestWeighted(long weightLimit) {
4574
        var limiter = isWeighted()
1 ✔
4575
            ? new WeightLimiter<K, V>(weightLimit)
1 ✔
4576
            : new SizeLimiter<K, V>((int) Math.min(weightLimit, cache.size()), weightLimit);
1 ✔
4577
        return cache.evictionOrder(/* hottest= */ false, transformer, limiter);
1 ✔
4578
      }
4579
      @Override
4580
      public <T extends @Nullable Object> T coldest(
4581
          Function<Stream<CacheEntry<K, V>>, T> mappingFunction) {
4582
        requireNonNull(mappingFunction);
1 ✔
4583
        return cache.evictionOrder(/* hottest= */ false, transformer, mappingFunction);
1 ✔
4584
      }
4585
      @Override public Map<K, V> hottest(int limit) {
4586
        int expectedSize = Math.min(limit, cache.size());
1 ✔
4587
        var limiter = new SizeLimiter<K, V>(expectedSize, limit);
1 ✔
4588
        return cache.evictionOrder(/* hottest= */ true, transformer, limiter);
1 ✔
4589
      }
4590
      @Override public Map<K, V> hottestWeighted(long weightLimit) {
4591
        var limiter = isWeighted()
1 ✔
4592
            ? new WeightLimiter<K, V>(weightLimit)
1 ✔
4593
            : new SizeLimiter<K, V>((int) Math.min(weightLimit, cache.size()), weightLimit);
1 ✔
4594
        return cache.evictionOrder(/* hottest= */ true, transformer, limiter);
1 ✔
4595
      }
4596
      @Override
4597
      public <T extends @Nullable Object> T hottest(
4598
          Function<Stream<CacheEntry<K, V>>, T> mappingFunction) {
4599
        requireNonNull(mappingFunction);
1 ✔
4600
        return cache.evictionOrder(/* hottest= */ true, transformer, mappingFunction);
1 ✔
4601
      }
4602
    }
4603

4604
    @SuppressWarnings("PreferJavaTimeOverload")
4605
    final class BoundedExpireAfterAccess implements FixedExpiration<K, V> {
1 ✔
4606
      @Override public OptionalLong ageOf(K key, TimeUnit unit) {
4607
        requireNonNull(key);
1 ✔
4608
        requireNonNull(unit);
1 ✔
4609
        Object lookupKey = cache.nodeFactory.newLookupKey(key);
1 ✔
4610
        Node<K, V> node = cache.data.get(lookupKey);
1 ✔
4611
        if (node == null) {
1 ✔
4612
          return OptionalLong.empty();
1 ✔
4613
        }
4614
        long now = cache.expirationTicker().read();
1 ✔
4615
        boolean expired = cache.hasExpired(node, now);
1 ✔
4616
        V value = node.getValue();
1 ✔
4617
        if ((value == null) || expired || cache.isComputingAsync(value)) {
1 ✔
4618
          return OptionalLong.empty();
1 ✔
4619
        }
4620
        long age = now - node.getAccessTime();
1 ✔
4621
        return (age < 0)
1 ✔
4622
            ? OptionalLong.empty()
1 ✔
4623
            : OptionalLong.of(unit.convert(age, TimeUnit.NANOSECONDS));
1 ✔
4624
      }
4625
      @Override public long getExpiresAfter(TimeUnit unit) {
4626
        return unit.convert(cache.expiresAfterAccessNanos(), TimeUnit.NANOSECONDS);
1 ✔
4627
      }
4628
      @Override public void setExpiresAfter(long duration, TimeUnit unit) {
4629
        requireArgument(duration >= 0, "duration cannot be negative: %s %s", duration, unit);
1 ✔
4630
        cache.setExpiresAfterAccessNanos(unit.toNanos(duration));
1 ✔
4631
        cache.scheduleAfterWrite();
1 ✔
4632
      }
1 ✔
4633
      @Override public Map<K, V> oldest(int limit) {
4634
        return oldest(new SizeLimiter<>(Math.min(limit, cache.size()), limit));
1 ✔
4635
      }
4636
      @Override public <T extends @Nullable Object> T oldest(
4637
          Function<Stream<CacheEntry<K, V>>, T> mappingFunction) {
4638
        return cache.expireAfterAccessOrder(/* oldest= */ true, transformer, mappingFunction);
1 ✔
4639
      }
4640
      @Override public Map<K, V> youngest(int limit) {
4641
        return youngest(new SizeLimiter<>(Math.min(limit, cache.size()), limit));
1 ✔
4642
      }
4643
      @Override public <T extends @Nullable Object> T youngest(
4644
          Function<Stream<CacheEntry<K, V>>, T> mappingFunction) {
4645
        return cache.expireAfterAccessOrder(/* oldest= */ false, transformer, mappingFunction);
1 ✔
4646
      }
4647
    }
4648

4649
    @SuppressWarnings("PreferJavaTimeOverload")
4650
    final class BoundedExpireAfterWrite implements FixedExpiration<K, V> {
1 ✔
4651
      @Override public OptionalLong ageOf(K key, TimeUnit unit) {
4652
        requireNonNull(key);
1 ✔
4653
        requireNonNull(unit);
1 ✔
4654
        Object lookupKey = cache.nodeFactory.newLookupKey(key);
1 ✔
4655
        Node<K, V> node = cache.data.get(lookupKey);
1 ✔
4656
        if (node == null) {
1 ✔
4657
          return OptionalLong.empty();
1 ✔
4658
        }
4659
        long now = cache.expirationTicker().read();
1 ✔
4660
        boolean expired = cache.hasExpired(node, now);
1 ✔
4661
        V value = node.getValue();
1 ✔
4662
        if ((value == null) || expired || cache.isComputingAsync(value)) {
1 ✔
4663
          return OptionalLong.empty();
1 ✔
4664
        }
4665
        long age = toWriteTime(now) - writeTimeOf(node);
1 ✔
4666
        return (age < 0)
1 ✔
4667
            ? OptionalLong.empty()
1 ✔
4668
            : OptionalLong.of(unit.convert(age, TimeUnit.NANOSECONDS));
1 ✔
4669
      }
4670
      @Override public long getExpiresAfter(TimeUnit unit) {
4671
        return unit.convert(cache.expiresAfterWriteNanos(), TimeUnit.NANOSECONDS);
1 ✔
4672
      }
4673
      @Override public void setExpiresAfter(long duration, TimeUnit unit) {
4674
        requireArgument(duration >= 0, "duration cannot be negative: %s %s", duration, unit);
1 ✔
4675
        cache.setExpiresAfterWriteNanos(unit.toNanos(duration));
1 ✔
4676
        cache.scheduleAfterWrite();
1 ✔
4677
      }
1 ✔
4678
      @Override public Map<K, V> oldest(int limit) {
4679
        return oldest(new SizeLimiter<>(Math.min(limit, cache.size()), limit));
1 ✔
4680
      }
4681
      @SuppressWarnings("GuardedByChecker")
4682
      @Override public <T extends @Nullable Object> T oldest(
4683
          Function<Stream<CacheEntry<K, V>>, T> mappingFunction) {
4684
        return cache.snapshot(cache.writeOrderDeque(), transformer, mappingFunction);
1 ✔
4685
      }
4686
      @Override public Map<K, V> youngest(int limit) {
4687
        return youngest(new SizeLimiter<>(Math.min(limit, cache.size()), limit));
1 ✔
4688
      }
4689
      @SuppressWarnings("GuardedByChecker")
4690
      @Override public <T extends @Nullable Object> T youngest(
4691
          Function<Stream<CacheEntry<K, V>>, T> mappingFunction) {
4692
        return cache.snapshot(cache.writeOrderDeque()::descendingIterator,
1 ✔
4693
            transformer, mappingFunction);
4694
      }
4695
    }
4696

4697
    @SuppressWarnings("PreferJavaTimeOverload")
4698
    final class BoundedVarExpiration implements VarExpiration<K, V> {
1 ✔
4699
      @Override public OptionalLong getExpiresAfter(K key, TimeUnit unit) {
4700
        requireNonNull(key);
1 ✔
4701
        requireNonNull(unit);
1 ✔
4702
        Object lookupKey = cache.nodeFactory.newLookupKey(key);
1 ✔
4703
        Node<K, V> node = cache.data.get(lookupKey);
1 ✔
4704
        if (node == null) {
1 ✔
4705
          return OptionalLong.empty();
1 ✔
4706
        }
4707
        long now = cache.expirationTicker().read();
1 ✔
4708
        boolean expired = cache.hasExpired(node, now);
1 ✔
4709
        V value = node.getValue();
1 ✔
4710
        if ((value == null) || expired || cache.isComputingAsync(value)) {
1 ✔
4711
          return OptionalLong.empty();
1 ✔
4712
        }
4713
        long duration = node.getVariableTime() - now;
1 ✔
4714
        return (cache.isAsync && (duration > MAXIMUM_EXPIRY))
1 ✔
4715
            ? OptionalLong.empty()
1 ✔
4716
            : OptionalLong.of(unit.convert(duration, TimeUnit.NANOSECONDS));
1 ✔
4717
      }
4718
      @Override public void setExpiresAfter(K key, long duration, TimeUnit unit) {
4719
        requireNonNull(key);
1 ✔
4720
        requireNonNull(unit);
1 ✔
4721
        requireArgument(duration >= 0, "duration cannot be negative: %s %s", duration, unit);
1 ✔
4722
        Object lookupKey = cache.nodeFactory.newLookupKey(key);
1 ✔
4723
        Node<K, V> node = cache.data.get(lookupKey);
1 ✔
4724
        if (node != null) {
1 ✔
4725
          long durationNanos = TimeUnit.NANOSECONDS.convert(duration, unit);
1 ✔
4726
          synchronized (node) {
1 ✔
4727
            long now = cache.expirationTicker().read();
1 ✔
4728
            V value = node.getValue();
1 ✔
4729
            if ((value == null) || cache.isComputingAsync(value)
1 ✔
4730
                || cache.hasExpired(node, now)) {
1 ✔
4731
              return;
1 ✔
4732
            }
4733
            node.setVariableTime(now + Math.min(durationNanos, MAXIMUM_EXPIRY));
1 ✔
4734
          }
1 ✔
4735
          cache.afterWrite(cache.new UpdateTask(node, 0, Access.QUIET));
1 ✔
4736
        }
4737
      }
1 ✔
4738
      @Override public @Nullable V put(K key, V value, long duration, TimeUnit unit) {
4739
        requireNonNull(unit);
1 ✔
4740
        requireNonNull(value);
1 ✔
4741
        requireArgument(duration >= 0, "duration cannot be negative: %s %s", duration, unit);
1 ✔
4742
        return cache.isAsync
1 ✔
4743
            ? putAsync(key, value, duration, unit)
1 ✔
4744
            : putSync(key, value, duration, unit, /* onlyIfAbsent= */ false);
1 ✔
4745
      }
4746
      @Override public @Nullable V putIfAbsent(K key, V value, long duration, TimeUnit unit) {
4747
        requireNonNull(unit);
1 ✔
4748
        requireNonNull(value);
1 ✔
4749
        requireArgument(duration >= 0, "duration cannot be negative: %s %s", duration, unit);
1 ✔
4750
        return cache.isAsync
1 ✔
4751
            ? putIfAbsentAsync(key, value, duration, unit)
1 ✔
4752
            : putSync(key, value, duration, unit, /* onlyIfAbsent= */ true);
1 ✔
4753
      }
4754
      @Nullable V putSync(K key, V value, long duration, TimeUnit unit, boolean onlyIfAbsent) {
4755
        var expiry = new FixedExpireAfterWrite<K, V>(duration, unit);
1 ✔
4756
        return cache.put(key, value, expiry, onlyIfAbsent);
1 ✔
4757
      }
4758
      @SuppressWarnings("unchecked")
4759
      @Nullable V putIfAbsentAsync(K key, V value, long duration, TimeUnit unit) {
4760
        // Mirrors LocalAsyncCache.AsMapView#putIfAbsent(key, value), but always probes quietly:
4761
        // the fixed duration must not be recomputed by the user's Expiry on a present entry
4762
        var expiry = (Expiry<K, V>) new AsyncExpiry<>(new FixedExpireAfterWrite<>(duration, unit));
1 ✔
4763
        var asyncValue = (V) CompletableFuture.completedFuture(value);
1 ✔
4764

4765
        for (;;) {
4766
          var priorFuture = (CompletableFuture<V>) cache.getIfPresentQuietly(key);
1 ✔
4767
          if (priorFuture != null) {
1 ✔
4768
            if (!priorFuture.isDone()) {
1 ✔
4769
              Async.getWhenSuccessful(priorFuture);
1 ✔
4770
              continue;
1 ✔
4771
            }
4772

4773
            V prior = Async.getWhenSuccessful(priorFuture);
1 ✔
4774
            if (prior != null) {
1 ✔
4775
              return prior;
1 ✔
4776
            }
4777
          }
4778

4779
          boolean[] added = { false };
1 ✔
4780
          var hints = new LocalCache.RemapHints();
1 ✔
4781
          var computed = (CompletableFuture<V>) cache.compute(key, (K k, @Nullable V oldValue) -> {
1 ✔
4782
            var oldValueFuture = (CompletableFuture<V>) oldValue;
1 ✔
4783
            added[0] = (oldValueFuture == null)
1 ✔
4784
                || (oldValueFuture.isDone() && (Async.getIfReady(oldValueFuture) == null));
1 ✔
4785
            if (added[0]) {
1 ✔
4786
              return asyncValue;
1 ✔
4787
            }
4788
            hints.preserveEntry();
1 ✔
4789
            return oldValue;
1 ✔
4790
          }, expiry, /* recordLoad= */ false, /* recordLoadFailure= */ false, hints);
4791

4792
          if (added[0]) {
1 ✔
4793
            return null;
1 ✔
4794
          } else {
4795
            V prior = Async.getWhenSuccessful(computed);
1 ✔
4796
            if (prior != null) {
1 ✔
4797
              return prior;
1 ✔
4798
            }
4799
          }
4800
        }
1 ✔
4801
      }
4802
      @SuppressWarnings("unchecked")
4803
      @Nullable V putAsync(K key, V value, long duration, TimeUnit unit) {
4804
        var expiry = (Expiry<K, V>) new AsyncExpiry<>(new FixedExpireAfterWrite<>(duration, unit));
1 ✔
4805
        var asyncValue = (V) CompletableFuture.completedFuture(value);
1 ✔
4806

4807
        var oldValueFuture = (CompletableFuture<V>) cache.put(
1 ✔
4808
            key, asyncValue, expiry, /* onlyIfAbsent= */ false);
4809
        return Async.getWhenSuccessful(oldValueFuture);
1 ✔
4810
      }
4811
      @Override public @Nullable V compute(K key,
4812
          BiFunction<? super K, ? super @Nullable V, ? extends @Nullable V> remappingFunction,
4813
          Duration duration) {
4814
        requireNonNull(key);
1 ✔
4815
        requireNonNull(duration);
1 ✔
4816
        requireNonNull(remappingFunction);
1 ✔
4817
        requireArgument(!duration.isNegative(), "duration cannot be negative: %s", duration);
1 ✔
4818
        var expiry = new FixedExpireAfterWrite<K, V>(
1 ✔
4819
            toNanosSaturated(duration), TimeUnit.NANOSECONDS);
1 ✔
4820

4821
        return cache.isAsync
1 ✔
4822
            ? computeAsync(key, remappingFunction, expiry)
1 ✔
4823
            : cache.compute(key, remappingFunction, expiry,
1 ✔
4824
                /* recordLoad= */ true, /* recordLoadFailure= */ true);
4825
      }
4826
      @Nullable V computeAsync(K key,
4827
          BiFunction<? super K, ? super @Nullable V, ? extends @Nullable V> remappingFunction,
4828
          Expiry<? super K, ? super V> expiry) {
4829
        // Keep in sync with LocalAsyncCache.AsMapView#compute(key, remappingFunction)
4830
        @SuppressWarnings("unchecked")
4831
        var delegate = (LocalCache<K, CompletableFuture<V>>) cache;
1 ✔
4832

4833
        @SuppressWarnings({"rawtypes", "unchecked", "Varifier"})
4834
        @Nullable V[] newValue = (@Nullable V[]) new Object[1];
1 ✔
4835
        for (;;) {
4836
          Async.getWhenSuccessful(delegate.getIfPresentQuietly(key));
1 ✔
4837

4838
          var hints = new LocalCache.RemapHints();
1 ✔
4839
          CompletableFuture<V> valueFuture = delegate.compute(
1 ✔
4840
              key, (K k, @Nullable CompletableFuture<V> oldValueFuture) -> {
4841
                if ((oldValueFuture != null) && !oldValueFuture.isDone()) {
1 ✔
4842
                  hints.preserveEntry();
1 ✔
4843
                  return oldValueFuture;
1 ✔
4844
                }
4845

4846
                V oldValue = Async.getIfReady(oldValueFuture);
1 ✔
4847
                BiFunction<? super K, ? super @Nullable V, ? extends @Nullable V> function =
1 ✔
4848
                    delegate.statsAware(remappingFunction,
1 ✔
4849
                        /* recordLoad= */ true, /* recordLoadFailure= */ true);
4850
                newValue[0] = function.apply(key, oldValue);
1 ✔
4851
                return (newValue[0] == null) ? null
1 ✔
4852
                    : CompletableFuture.completedFuture(newValue[0]);
1 ✔
4853
              }, new AsyncExpiry<>(expiry), /* recordLoad= */ false,
4854
              /* recordLoadFailure= */ false, hints);
4855

4856
          if (newValue[0] != null) {
1 ✔
4857
            return newValue[0];
1 ✔
4858
          } else if (valueFuture == null) {
1 ✔
4859
            return null;
1 ✔
4860
          }
4861
        }
1 ✔
4862
      }
4863
      @Override public Map<K, V> oldest(int limit) {
4864
        return oldest(new SizeLimiter<>(Math.min(limit, cache.size()), limit));
1 ✔
4865
      }
4866
      @Override public <T extends @Nullable Object> T oldest(
4867
          Function<Stream<CacheEntry<K, V>>, T> mappingFunction) {
4868
        return cache.snapshot(cache.timerWheel(), transformer, mappingFunction);
1 ✔
4869
      }
4870
      @Override public Map<K, V> youngest(int limit) {
4871
        return youngest(new SizeLimiter<>(Math.min(limit, cache.size()), limit));
1 ✔
4872
      }
4873
      @Override public <T extends @Nullable Object> T youngest(
4874
          Function<Stream<CacheEntry<K, V>>, T> mappingFunction) {
4875
        return cache.snapshot(cache.timerWheel()::descendingIterator, transformer, mappingFunction);
1 ✔
4876
      }
4877
    }
4878

4879
    static final class FixedExpireAfterWrite<K, V> implements Expiry<K, V> {
4880
      final long duration;
4881
      final TimeUnit unit;
4882

4883
      FixedExpireAfterWrite(long duration, TimeUnit unit) {
1 ✔
4884
        this.duration = duration;
1 ✔
4885
        this.unit = unit;
1 ✔
4886
      }
1 ✔
4887
      @Override public long expireAfterCreate(K key, V value, long currentTime) {
4888
        return unit.toNanos(duration);
1 ✔
4889
      }
4890
      @Override public long expireAfterUpdate(
4891
          K key, V value, long currentTime, long currentDuration) {
4892
        return unit.toNanos(duration);
1 ✔
4893
      }
4894
      @CanIgnoreReturnValue
4895
      @Override public long expireAfterRead(
4896
          K key, V value, long currentTime, long currentDuration) {
4897
        return currentDuration;
1 ✔
4898
      }
4899
    }
4900

4901
    @SuppressWarnings("PreferJavaTimeOverload")
4902
    final class BoundedRefreshAfterWrite implements FixedRefresh<K, V> {
1 ✔
4903
      @Override public OptionalLong ageOf(K key, TimeUnit unit) {
4904
        requireNonNull(key);
1 ✔
4905
        requireNonNull(unit);
1 ✔
4906
        Object lookupKey = cache.nodeFactory.newLookupKey(key);
1 ✔
4907
        Node<K, V> node = cache.data.get(lookupKey);
1 ✔
4908
        if (node == null) {
1 ✔
4909
          return OptionalLong.empty();
1 ✔
4910
        }
4911
        long now = cache.expirationTicker().read();
1 ✔
4912
        boolean expired = cache.hasExpired(node, now);
1 ✔
4913
        V value = node.getValue();
1 ✔
4914
        if ((value == null) || expired || cache.isComputingAsync(value)) {
1 ✔
4915
          return OptionalLong.empty();
1 ✔
4916
        }
4917
        long age = toWriteTime(now) - writeTimeOf(node);
1 ✔
4918
        return (age < 0)
1 ✔
4919
            ? OptionalLong.empty()
1 ✔
4920
            : OptionalLong.of(unit.convert(age, TimeUnit.NANOSECONDS));
1 ✔
4921
      }
4922
      @Override public long getRefreshesAfter(TimeUnit unit) {
4923
        return unit.convert(cache.refreshAfterWriteNanos(), TimeUnit.NANOSECONDS);
1 ✔
4924
      }
4925
      @Override public void setRefreshesAfter(long duration, TimeUnit unit) {
4926
        requireNonNull(unit);
1 ✔
4927
        requireArgument(duration > 0, "duration must be positive: %s %s", duration, unit);
1 ✔
4928
        cache.setRefreshAfterWriteNanos(unit.toNanos(duration));
1 ✔
4929
        cache.scheduleAfterWrite();
1 ✔
4930
      }
1 ✔
4931
    }
4932
  }
4933

4934
  /* --------------- Loading Cache --------------- */
4935

4936
  static final class BoundedLocalLoadingCache<K, V>
4937
      extends BoundedLocalManualCache<K, V> implements LocalLoadingCache<K, V> {
4938
    private static final long serialVersionUID = 1;
4939

4940
    final Function<K, @Nullable V> mappingFunction;
4941
    final @Nullable Function<Set<? extends K>, Map<K, V>> bulkMappingFunction;
4942

4943
    BoundedLocalLoadingCache(Caffeine<K, V> builder, CacheLoader<? super K, V> loader) {
4944
      super(builder, loader);
1 ✔
4945
      requireNonNull(loader);
1 ✔
4946
      mappingFunction = newMappingFunction(loader);
1 ✔
4947
      bulkMappingFunction = newBulkMappingFunction(loader);
1 ✔
4948
    }
1 ✔
4949

4950
    @Override
4951
    @SuppressWarnings({"DataFlowIssue", "NullAway"})
4952
    public AsyncCacheLoader<? super K, V> cacheLoader() {
4953
      return cache.cacheLoader;
1 ✔
4954
    }
4955

4956
    @Override
4957
    public Function<K, @Nullable V> mappingFunction() {
4958
      return mappingFunction;
1 ✔
4959
    }
4960

4961
    @Override
4962
    public @Nullable Function<Set<? extends K>, Map<K, V>> bulkMappingFunction() {
4963
      return bulkMappingFunction;
1 ✔
4964
    }
4965

4966
    private void readObject(ObjectInputStream stream) throws InvalidObjectException {
4967
      throw new InvalidObjectException("Proxy required");
1 ✔
4968
    }
4969

4970
    private Object writeReplace() {
4971
      return makeSerializationProxy(cache);
1 ✔
4972
    }
4973
  }
4974

4975
  /* --------------- Async Cache --------------- */
4976

4977
  static final class BoundedLocalAsyncCache<K, V> implements LocalAsyncCache<K, V>, Serializable {
4978
    private static final long serialVersionUID = 1;
4979

4980
    final BoundedLocalCache<K, CompletableFuture<V>> cache;
4981
    final boolean isWeighted;
4982

4983
    @Nullable ConcurrentMap<K, CompletableFuture<V>> mapView;
4984
    @Nullable CacheView<K, V> cacheView;
4985
    @Nullable Policy<K, V> policy;
4986

4987
    @SuppressWarnings("unchecked")
4988
    BoundedLocalAsyncCache(Caffeine<K, V> builder) {
1 ✔
4989
      cache = (BoundedLocalCache<K, CompletableFuture<V>>) LocalCacheFactory
1 ✔
4990
          .newBoundedLocalCache(builder, /* cacheLoader= */ null, /* isAsync= */ true);
1 ✔
4991
      isWeighted = builder.isWeighted();
1 ✔
4992
    }
1 ✔
4993

4994
    @Override
4995
    public BoundedLocalCache<K, CompletableFuture<V>> cache() {
4996
      return cache;
1 ✔
4997
    }
4998

4999
    @Override
5000
    public ConcurrentMap<K, CompletableFuture<V>> asMap() {
5001
      return (mapView == null) ? (mapView = new AsyncAsMapView<>(this)) : mapView;
1 ✔
5002
    }
5003

5004
    @Override
5005
    public Cache<K, V> synchronous() {
5006
      return (cacheView == null) ? (cacheView = new CacheView<>(this)) : cacheView;
1 ✔
5007
    }
5008

5009
    @Override
5010
    public Policy<K, V> policy() {
5011
      if (policy == null) {
1 ✔
5012
        @SuppressWarnings("unchecked")
5013
        var castCache = (BoundedLocalCache<K, V>) cache;
1 ✔
5014
        Function<CompletableFuture<V>, @Nullable V> transformer = Async::getIfReady;
1 ✔
5015
        @SuppressWarnings("unchecked")
5016
        var castTransformer = (Function<@Nullable V, @Nullable V>) transformer;
1 ✔
5017
        policy = new BoundedPolicy<>(castCache, castTransformer, isWeighted);
1 ✔
5018
      }
5019
      return policy;
1 ✔
5020
    }
5021

5022
    private void readObject(ObjectInputStream stream) throws InvalidObjectException {
5023
      throw new InvalidObjectException("Proxy required");
1 ✔
5024
    }
5025

5026
    private Object writeReplace() {
5027
      return makeSerializationProxy(cache);
1 ✔
5028
    }
5029
  }
5030

5031
  /* --------------- Async Loading Cache --------------- */
5032

5033
  static final class BoundedLocalAsyncLoadingCache<K, V>
5034
      extends LocalAsyncLoadingCache<K, V> implements Serializable {
5035
    private static final long serialVersionUID = 1;
5036

5037
    final BoundedLocalCache<K, CompletableFuture<V>> cache;
5038
    final boolean isWeighted;
5039

5040
    @Nullable ConcurrentMap<K, CompletableFuture<V>> mapView;
5041
    @Nullable Policy<K, V> policy;
5042

5043
    @SuppressWarnings("unchecked")
5044
    BoundedLocalAsyncLoadingCache(Caffeine<K, V> builder, AsyncCacheLoader<? super K, V> loader) {
5045
      super(loader);
1 ✔
5046
      isWeighted = builder.isWeighted();
1 ✔
5047
      cache = (BoundedLocalCache<K, CompletableFuture<V>>) LocalCacheFactory
1 ✔
5048
          .newBoundedLocalCache(builder, loader, /* isAsync= */ true);
1 ✔
5049
    }
1 ✔
5050

5051
    @Override
5052
    public BoundedLocalCache<K, CompletableFuture<V>> cache() {
5053
      return cache;
1 ✔
5054
    }
5055

5056
    @Override
5057
    public ConcurrentMap<K, CompletableFuture<V>> asMap() {
5058
      return (mapView == null) ? (mapView = new AsyncAsMapView<>(this)) : mapView;
1 ✔
5059
    }
5060

5061
    @Override
5062
    public Policy<K, V> policy() {
5063
      if (policy == null) {
1 ✔
5064
        @SuppressWarnings("unchecked")
5065
        var castCache = (BoundedLocalCache<K, V>) cache;
1 ✔
5066
        Function<CompletableFuture<V>, @Nullable V> transformer = Async::getIfReady;
1 ✔
5067
        @SuppressWarnings("unchecked")
5068
        var castTransformer = (Function<@Nullable V, @Nullable V>) transformer;
1 ✔
5069
        policy = new BoundedPolicy<>(castCache, castTransformer, isWeighted);
1 ✔
5070
      }
5071
      return policy;
1 ✔
5072
    }
5073

5074
    private void readObject(ObjectInputStream stream) throws InvalidObjectException {
5075
      throw new InvalidObjectException("Proxy required");
1 ✔
5076
    }
5077

5078
    private Object writeReplace() {
5079
      return makeSerializationProxy(cache);
1 ✔
5080
    }
5081
  }
5082
}
5083

5084
/** The namespace for field padding through inheritance. */
5085
@SuppressWarnings({"IdentifierName", "MultiVariableDeclaration"})
5086
final class BLCHeader {
5087

5088
  private BLCHeader() {}
5089

5090
  @SuppressWarnings("unused")
5091
  static class PadDrainStatus {
1 ✔
5092
    byte p000, p001, p002, p003, p004, p005, p006, p007;
5093
    byte p008, p009, p010, p011, p012, p013, p014, p015;
5094
    byte p016, p017, p018, p019, p020, p021, p022, p023;
5095
    byte p024, p025, p026, p027, p028, p029, p030, p031;
5096
    byte p032, p033, p034, p035, p036, p037, p038, p039;
5097
    byte p040, p041, p042, p043, p044, p045, p046, p047;
5098
    byte p048, p049, p050, p051, p052, p053, p054, p055;
5099
    byte p056, p057, p058, p059, p060, p061, p062, p063;
5100
    byte p064, p065, p066, p067, p068, p069, p070, p071;
5101
    byte p072, p073, p074, p075, p076, p077, p078, p079;
5102
    byte p080, p081, p082, p083, p084, p085, p086, p087;
5103
    byte p088, p089, p090, p091, p092, p093, p094, p095;
5104
    byte p096, p097, p098, p099, p100, p101, p102, p103;
5105
    byte p104, p105, p106, p107, p108, p109, p110, p111;
5106
    byte p112, p113, p114, p115, p116, p117, p118, p119;
5107
  }
5108

5109
  /** Enforces a memory layout to avoid false sharing by padding the drain status. */
5110
  abstract static class DrainStatusRef extends PadDrainStatus {
1 ✔
5111
    static final VarHandle DRAIN_STATUS = fieldVarHandle(MethodHandles.lookup(),
1 ✔
5112
        "drainStatus", VarHandle.class, DrainStatusRef.class, int.class);
5113

5114
    /** A drain is not taking place. */
5115
    static final int IDLE = 0;
5116
    /** A drain is required due to a pending write modification. */
5117
    static final int REQUIRED = 1;
5118
    /** A drain is in progress and will transition to idle. */
5119
    static final int PROCESSING_TO_IDLE = 2;
5120
    /** A drain is in progress and will transition to required. */
5121
    static final int PROCESSING_TO_REQUIRED = 3;
5122

5123
    /** The draining status of the buffers. */
5124
    volatile int drainStatus = IDLE;
1 ✔
5125

5126
    /**
5127
     * Returns whether maintenance work is needed.
5128
     *
5129
     * @param delayable if draining the read buffer can be delayed
5130
     */
5131
    @SuppressWarnings("StatementSwitchToExpressionSwitch")
5132
    boolean shouldDrainBuffers(boolean delayable) {
5133
      switch (drainStatusOpaque()) {
1 ✔
5134
        case IDLE:
5135
          return !delayable;
1 ✔
5136
        case REQUIRED:
5137
          return true;
1 ✔
5138
        case PROCESSING_TO_IDLE:
5139
        case PROCESSING_TO_REQUIRED:
5140
          return false;
1 ✔
5141
        default:
5142
          throw new IllegalStateException("Invalid drain status: " + drainStatus);
1 ✔
5143
      }
5144
    }
5145

5146
    int drainStatusOpaque() {
5147
      return (int) DRAIN_STATUS.getOpaque(this);
1 ✔
5148
    }
5149

5150
    int drainStatusAcquire() {
5151
      return (int) DRAIN_STATUS.getAcquire(this);
1 ✔
5152
    }
5153

5154
    void setDrainStatusOpaque(int drainStatus) {
5155
      DRAIN_STATUS.setOpaque(this, drainStatus);
1 ✔
5156
    }
1 ✔
5157

5158
    void setDrainStatusRelease(int drainStatus) {
5159
      DRAIN_STATUS.setRelease(this, drainStatus);
1 ✔
5160
    }
1 ✔
5161

5162
    boolean casDrainStatus(int expect, int update) {
5163
      return DRAIN_STATUS.compareAndSet(this, expect, update);
1 ✔
5164
    }
5165
  }
5166
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc