• 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.88
/jcache/src/main/java/com/github/benmanes/caffeine/jcache/CacheProxy.java
1
/*
2
 * Copyright 2015 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.jcache;
17

18
import static java.util.Objects.requireNonNull;
19
import static java.util.Objects.requireNonNullElse;
20
import static java.util.stream.Collectors.toMap;
21
import static java.util.stream.Collectors.toSet;
22
import static java.util.stream.Collectors.toUnmodifiableList;
23

24
import java.lang.System.Logger;
25
import java.lang.System.Logger.Level;
26
import java.util.ArrayList;
27
import java.util.HashMap;
28
import java.util.Iterator;
29
import java.util.LinkedHashSet;
30
import java.util.List;
31
import java.util.Map;
32
import java.util.NoSuchElementException;
33
import java.util.Objects;
34
import java.util.Optional;
35
import java.util.Set;
36
import java.util.concurrent.CompletableFuture;
37
import java.util.concurrent.ConcurrentHashMap;
38
import java.util.concurrent.ExecutionException;
39
import java.util.concurrent.Executor;
40
import java.util.concurrent.ExecutorService;
41
import java.util.concurrent.TimeUnit;
42
import java.util.concurrent.TimeoutException;
43
import java.util.function.BiFunction;
44
import java.util.function.Consumer;
45
import java.util.function.Supplier;
46

47
import javax.cache.Cache;
48
import javax.cache.CacheException;
49
import javax.cache.CacheManager;
50
import javax.cache.configuration.CacheEntryListenerConfiguration;
51
import javax.cache.configuration.Configuration;
52
import javax.cache.event.CacheEntryListenerException;
53
import javax.cache.expiry.Duration;
54
import javax.cache.expiry.ExpiryPolicy;
55
import javax.cache.integration.CacheLoader;
56
import javax.cache.integration.CacheLoaderException;
57
import javax.cache.integration.CacheWriter;
58
import javax.cache.integration.CacheWriterException;
59
import javax.cache.integration.CompletionListener;
60
import javax.cache.processor.EntryProcessor;
61
import javax.cache.processor.EntryProcessorException;
62
import javax.cache.processor.EntryProcessorResult;
63

64
import org.jspecify.annotations.NonNull;
65
import org.jspecify.annotations.Nullable;
66

67
import com.github.benmanes.caffeine.cache.Ticker;
68
import com.github.benmanes.caffeine.jcache.configuration.CaffeineConfiguration;
69
import com.github.benmanes.caffeine.jcache.copy.Copier;
70
import com.github.benmanes.caffeine.jcache.event.EventDispatcher;
71
import com.github.benmanes.caffeine.jcache.event.Registration;
72
import com.github.benmanes.caffeine.jcache.integration.DisabledCacheWriter;
73
import com.github.benmanes.caffeine.jcache.management.JCacheMXBean;
74
import com.github.benmanes.caffeine.jcache.management.JCacheStatisticsMXBean;
75
import com.github.benmanes.caffeine.jcache.management.JmxRegistration;
76
import com.github.benmanes.caffeine.jcache.management.JmxRegistration.MBeanType;
77
import com.github.benmanes.caffeine.jcache.processor.Action;
78
import com.github.benmanes.caffeine.jcache.processor.EntryProcessorEntry;
79
import com.google.errorprone.annotations.CanIgnoreReturnValue;
80
import com.google.errorprone.annotations.Var;
81

82
/**
83
 * An implementation of JSR-107 {@link Cache} backed by a Caffeine cache.
84
 *
85
 * @author ben.manes@gmail.com (Ben Manes)
86
 */
87
@SuppressWarnings("OptionalUsedAsFieldOrParameterType")
88
public class CacheProxy<K, V> implements Cache<K, V> {
89
  private static final Logger logger = System.getLogger(CacheProxy.class.getName());
1 ✔
90

91
  protected final com.github.benmanes.caffeine.cache.Cache<K, @Nullable Expirable<V>> cache;
92
  protected final Optional<CacheLoader<K, V>> cacheLoader;
93
  protected final Set<CompletableFuture<?>> inFlight;
94
  protected final JCacheStatisticsMXBean statistics;
95
  protected final EventDispatcher<K, V> dispatcher;
96
  protected final Executor executor;
97
  protected final Ticker ticker;
98

99
  private final CaffeineConfiguration<K, V> configuration;
100
  private final CacheManagerImpl cacheManager;
101
  private final CacheWriter<K, V> writer;
102
  private final JCacheMXBean cacheMxBean;
103
  private final ExpiryPolicy expiry;
104
  private final Copier copier;
105
  private final String name;
106

107
  private volatile boolean closed;
108

109
  @SuppressWarnings({"PMD.ExcessiveParameterList", "this-escape", "TooManyParameters"})
110
  public CacheProxy(String name, Executor executor, CacheManagerImpl cacheManager,
111
      CaffeineConfiguration<K, V> configuration,
112
      com.github.benmanes.caffeine.cache.Cache<K, @Nullable Expirable<V>> cache,
113
      EventDispatcher<K, V> dispatcher, Optional<CacheLoader<K, V>> cacheLoader,
114
      ExpiryPolicy expiry, Ticker ticker, JCacheStatisticsMXBean statistics, Copier copier) {
1 ✔
115
    var cacheWriter = configuration.getCacheWriter();
1 ✔
116
    this.writer = configuration.isWriteThrough()
1 ✔
117
        ? requireNonNull(cacheWriter, "The CacheWriter factory returned null")
1 ✔
118
        : requireNonNullElse(cacheWriter, DisabledCacheWriter.get());
1 ✔
119
    this.configuration = requireNonNull(configuration);
1 ✔
120
    this.cacheManager = requireNonNull(cacheManager);
1 ✔
121
    this.cacheLoader = requireNonNull(cacheLoader);
1 ✔
122
    this.inFlight = ConcurrentHashMap.newKeySet();
1 ✔
123
    this.dispatcher = requireNonNull(dispatcher);
1 ✔
124
    this.statistics = requireNonNull(statistics);
1 ✔
125
    this.cacheMxBean = new JCacheMXBean(this);
1 ✔
126
    this.executor = requireNonNull(executor);
1 ✔
127
    this.expiry = requireNonNull(expiry);
1 ✔
128
    this.ticker = requireNonNull(ticker);
1 ✔
129
    this.cache = requireNonNull(cache);
1 ✔
130
    this.copier = requireNonNull(copier);
1 ✔
131
    this.name = requireNonNull(name);
1 ✔
132
  }
1 ✔
133

134
  @Override
135
  public boolean containsKey(K key) {
136
    requireOperable();
1 ✔
137

138
    Expirable<V> expirable = cache.getIfPresent(key);
1 ✔
139
    if (expirable == null) {
1 ✔
140
      return false;
1 ✔
141
    }
142
    if (expirable.isEternal()) {
1 ✔
143
      return true;
1 ✔
144
    }
145
    long millis = currentTimeMillis();
1 ✔
146
    if (expirable.hasExpired(millis)) {
1 ✔
147
      dispatcher.beginComputation();
1 ✔
148
      try {
149
        cache.asMap().computeIfPresent(key, (k, e) -> {
1 ✔
150
          if ((e == expirable) && expirable.hasExpired(millis)) {
1 ✔
151
            dispatcher.publishExpired(this, key, expirable.get());
1 ✔
152
            statistics.recordEvictions(1L);
1 ✔
153
            return null;
1 ✔
154
          }
155
          return e;
1 ✔
156
        });
157
      } finally {
158
        dispatcher.endComputation();
1 ✔
159
      }
160
      dispatcher.awaitSynchronous();
1 ✔
161
      return false;
1 ✔
162
    }
163
    return true;
1 ✔
164
  }
165

166
  @Override
167
  public @Nullable V get(K key) {
168
    requireOperable();
1 ✔
169

170
    boolean statsEnabled = statistics.isEnabled();
1 ✔
171
    long start = statsEnabled ? ticker.read() : 0L;
1 ✔
172
    @Var Expirable<V> expirable = cache.getIfPresent(key);
1 ✔
173
    if (expirable == null) {
1 ✔
174
      statistics.recordMisses(1L);
1 ✔
175
      if (statsEnabled) {
1 ✔
176
        statistics.recordGetTime(ticker.read() - start);
1 ✔
177
      }
178
      return null;
1 ✔
179
    }
180

181
    long millis;
182
    if (expirable.isEternal()) {
1 ✔
183
      millis = 0L;
1 ✔
184
    } else {
185
      long now = ticker.read();
1 ✔
186
      millis = nanosToMillis(now);
1 ✔
187
      if (expirable.hasExpired(millis)) {
1 ✔
188
        Expirable<V> current;
189
        var expired = expirable;
1 ✔
190
        dispatcher.beginComputation();
1 ✔
191
        try {
192
          current = cache.asMap().computeIfPresent(key, (k, e) -> {
1 ✔
193
            if ((e == expired) && expired.hasExpired(millis)) {
1 ✔
194
              dispatcher.publishExpired(this, key, expired.get());
1 ✔
195
              statistics.recordEvictions(1L);
1 ✔
196
              return null;
1 ✔
197
            }
198
            return e;
1 ✔
199
          });
200
        } finally {
201
          dispatcher.endComputation();
1 ✔
202
        }
203
        if ((current == null) || current.hasExpired(millis)) {
1 !
204
          var listenerFailure = awaitSynchronousFailure();
1 ✔
205
          statistics.recordMisses(1L);
1 ✔
206
          if (statsEnabled) {
1 ✔
207
            statistics.recordGetTime(ticker.read() - start);
1 ✔
208
          }
209
          rethrowListenerFailure(listenerFailure);
1 ✔
210
          return null;
1 ✔
211
        }
212
        expirable = current;
1 ✔
213
      }
214
    }
215

216
    var duration = getAccessExpireTime();
1 ✔
217
    setAccessExpireTime(expirable, duration, millis);
1 ✔
218
    setVariableExpiration(key, duration);
1 ✔
219
    V value = copyOf(expirable.get());
1 ✔
220
    if (statsEnabled) {
1 ✔
221
      statistics.recordHits(1L);
1 ✔
222
      statistics.recordGetTime(ticker.read() - start);
1 ✔
223
    }
224
    return value;
1 ✔
225
  }
226

227
  @Override
228
  public Map<K, V> getAll(Set<? extends K> keys) {
229
    requireOperable();
1 ✔
230

231
    boolean statsEnabled = statistics.isEnabled();
1 ✔
232
    long now = statsEnabled ? ticker.read() : 0L;
1 ✔
233
    Map<K, V> result;
234
    try {
235
      Map<K, Expirable<V>> entries = getAndFilterExpiredEntries(keys);
1 ✔
236
      result = copyMap(entries);
1 ✔
237
      if (statsEnabled) {
1 ✔
238
        statistics.recordGetTime(ticker.read() - now);
1 ✔
239
      }
240
    } catch (Throwable t) {
1 ✔
241
      awaitAndSuppressFailure(t);
1 ✔
242
      throw t;
1 ✔
243
    }
1 ✔
244
    rethrowListenerFailure(awaitSynchronousFailure());
1 ✔
245
    return result;
1 ✔
246
  }
247

248
  /**
249
   * Returns all of the mappings present, expiring as required, and updates their access expiry
250
   * time.
251
   */
252
  protected Map<K, Expirable<V>> getAndFilterExpiredEntries(Set<? extends K> keys) {
253
    int[] expired = { 0 };
1 ✔
254
    long[] millis = { 0L };
1 ✔
255
    int requested = keys.size();
1 ✔
256
    var result = new HashMap<K, @NonNull Expirable<V>>(cache.getAllPresent(keys));
1 ✔
257
    result.entrySet().removeIf(entry -> {
1 ✔
258
      if (!entry.getValue().isEternal() && (millis[0] == 0L)) {
1 ✔
259
        millis[0] = currentTimeMillis();
1 ✔
260
      }
261
      if (!entry.getValue().isEternal() && entry.getValue().hasExpired(millis[0])) {
1 ✔
262
        dispatcher.beginComputation();
1 ✔
263
        Expirable<V> current;
264
        try {
265
          current = cache.asMap().computeIfPresent(entry.getKey(), (k, expirable) -> {
1 ✔
266
            if ((expirable == entry.getValue()) && expirable.hasExpired(millis[0])) {
1 ✔
267
              dispatcher.publishExpired(this, entry.getKey(), entry.getValue().get());
1 ✔
268
              expired[0]++;
1 ✔
269
              return null;
1 ✔
270
            }
271
            return expirable;
1 ✔
272
          });
273
        } finally {
274
          dispatcher.endComputation();
1 ✔
275
        }
276
        if ((current == null) || current.hasExpired(millis[0])) {
1 !
277
          return true;
1 ✔
278
        }
279
        entry.setValue(current);
1 ✔
280
      }
281
      var duration = getAccessExpireTime();
1 ✔
282
      setAccessExpireTime(entry.getValue(), duration, millis[0]);
1 ✔
283
      setVariableExpiration(entry.getKey(), duration);
1 ✔
284
      return false;
1 ✔
285
    });
286

287
    statistics.recordHits(result.size());
1 ✔
288
    statistics.recordMisses(requested - result.size());
1 ✔
289
    statistics.recordEvictions(expired[0]);
1 ✔
290
    return result;
1 ✔
291
  }
292

293
  @Override
294
  @SuppressWarnings({"CollectionUndefinedEquality", "FutureReturnValueIgnored"})
295
  public void loadAll(Set<? extends K> keys, boolean replaceExistingValues,
296
      @Nullable CompletionListener completionListener) {
297
    requireOperable();
1 ✔
298
    keys.forEach(Objects::requireNonNull);
1 ✔
299
    CompletionListener listener = (completionListener == null)
1 ✔
300
        ? NullCompletionListener.INSTANCE
1 ✔
301
        : completionListener;
1 ✔
302
    if (cacheLoader.isEmpty()) {
1 ✔
303
      listener.onCompletion();
1 ✔
304
      return;
1 ✔
305
    }
306

307
    var future = new CompletableFuture<@Nullable Void>();
1 ✔
308
    synchronized (configuration) {
1 ✔
309
      requireNotClosed();
1 ✔
310
      inFlight.add(future);
1 ✔
311
    }
1 ✔
312
    try {
313
      // The tracked future spans the notification so that close() awaits the listener's callback
314
      CompletableFuture
1 ✔
315
          .supplyAsync(() -> loadAllAndNotify(keys, replaceExistingValues, listener), executor)
1 ✔
316
          .thenCompose(chain -> chain)
1 ✔
317
          .whenComplete((r, e) -> {
1 ✔
318
            inFlight.remove(future);
1 ✔
319
            future.complete(null);
1 ✔
320
          });
1 ✔
321
    } catch (RuntimeException e) {
1 ✔
322
      inFlight.remove(future);
1 ✔
323
      future.complete(null);
1 ✔
324
      listener.onException(new CacheLoaderException(e));
1 ✔
325
    }
1 ✔
326
  }
1 ✔
327

328
  /** Performs the bulk load and returns a future that includes its completion notification. */
329
  private CompletableFuture<@Nullable Void> loadAllAndNotify(Set<? extends K> keys,
330
      boolean replaceExistingValues, CompletionListener listener) {
331
    @Var CacheLoaderException failure = null;
1 ✔
332
    try {
333
      if (replaceExistingValues) {
1 ✔
334
        loadAllAndReplaceExisting(keys);
1 ✔
335
      } else {
336
        loadAllAndKeepExisting(keys);
1 ✔
337
      }
338
    } catch (CacheLoaderException e) {
1 ✔
339
      failure = e;
1 ✔
340
    } catch (Throwable t) {
1 ✔
341
      restoreInterrupt(t);
1 ✔
342
      failure = new CacheLoaderException(t);
1 ✔
343
    }
1 ✔
344
    var loadFailure = failure;
1 ✔
345
    return dispatcher.chainSynchronous().<@Nullable Void>handle((listenerFailure, error) -> {
1 ✔
346
      var dispatchFailure = (error == null) ? listenerFailure : new CacheLoaderException(error);
1 ✔
347
      var outcome = suppress(loadFailure, dispatchFailure);
1 ✔
348
      if (outcome == null) {
1 ✔
349
        listener.onCompletion();
1 ✔
350
      } else {
351
        listener.onException(outcome);
1 ✔
352
      }
353
      return null;
1 ✔
354
    });
355
  }
356

357
  /** Performs the bulk load where the existing entries are replaced. */
358
  protected void loadAllAndReplaceExisting(Set<? extends K> keys) {
359
    Map<K, V> loaded = cacheLoader.orElseThrow().loadAll(keys);
1 ✔
360
    for (var entry : loaded.entrySet()) {
1 ✔
361
      if ((entry.getKey() != null) && (entry.getValue() != null)) {
1 ✔
362
        putNoCopyOrAwait(entry.getKey(), entry.getValue(), /* publishToWriter= */ false);
1 ✔
363
      }
364
    }
1 ✔
365
  }
1 ✔
366

367
  /** Performs the bulk load where the existing entries are retained. */
368
  @SuppressWarnings("ConstantValue")
369
  protected void loadAllAndKeepExisting(Set<? extends K> keys) {
370
    List<K> keysToLoad = keys.stream()
1 ✔
371
        .filter(key -> {
1 ✔
372
          var expirable = cache.policy().getIfPresentQuietly(key);
1 ✔
373
          return (expirable == null)
1 ✔
374
              || (!expirable.isEternal() && expirable.hasExpired(currentTimeMillis()));
1 ✔
375
        }).collect(toUnmodifiableList());
1 ✔
376
    Map<K, V> result = cacheLoader.orElseThrow().loadAll(keysToLoad);
1 ✔
377
    for (var entry : result.entrySet()) {
1 ✔
378
      if ((entry.getKey() != null) && (entry.getValue() != null)) {
1 ✔
379
        putIfAbsentNoAwait(entry.getKey(), entry.getValue(), /* publishToWriter= */ false);
1 ✔
380
      }
381
    }
1 ✔
382
  }
1 ✔
383

384
  @Override
385
  public void put(K key, V value) {
386
    requireOperable();
1 ✔
387

388
    boolean statsEnabled = statistics.isEnabled();
1 ✔
389
    long start = statsEnabled ? ticker.read() : 0L;
1 ✔
390
    var result = putNoCopyOrAwait(key, value, /* publishToWriter= */ true);
1 ✔
391
    var listenerFailure = awaitSynchronousFailure();
1 ✔
392
    if (statsEnabled && result.written) {
1 ✔
393
      statistics.recordPuts(1);
1 ✔
394
      statistics.recordPutTime(ticker.read() - start);
1 ✔
395
    }
396
    rethrowListenerFailure(listenerFailure);
1 ✔
397
  }
1 ✔
398

399
  @Override
400
  public @Nullable V getAndPut(K key, V value) {
401
    requireOperable();
1 ✔
402

403
    boolean statsEnabled = statistics.isEnabled();
1 ✔
404
    long start = statsEnabled ? ticker.read() : 0L;
1 ✔
405
    var result = putNoCopyOrAwait(key, value, /* publishToWriter= */ true);
1 ✔
406
    var listenerFailure = awaitSynchronousFailure();
1 ✔
407
    if (statsEnabled) {
1 ✔
408
      if (result.oldValue == null) {
1 ✔
409
        statistics.recordMisses(1L);
1 ✔
410
      } else {
411
        statistics.recordHits(1L);
1 ✔
412
      }
413
      long duration = ticker.read() - start;
1 ✔
414
      if (result.written) {
1 ✔
415
        statistics.recordPuts(1);
1 ✔
416
        statistics.recordPutTime(duration);
1 ✔
417
      }
418
      statistics.recordGetTime(duration);
1 ✔
419
    }
420
    rethrowListenerFailure(listenerFailure);
1 ✔
421
    return (result.oldValue == null) ? null : copyOf(result.oldValue);
1 ✔
422
  }
423

424
  /**
425
   * Associates the specified value with the specified key in the cache.
426
   *
427
   * @param key key with which the specified value is to be associated
428
   * @param value value to be associated with the specified key
429
   * @param publishToWriter if the writer should be notified
430
   * @return the oldValue and if the new value was stored
431
   */
432
  @CanIgnoreReturnValue
433
  protected PutResult<V> putNoCopyOrAwait(K key, V value, boolean publishToWriter) {
434
    requireNonNull(key);
1 ✔
435
    requireNonNull(value);
1 ✔
436
    var entry = new CopiedEntry<>(key, copyOf(key), value, copyOf(value));
1 ✔
437
    return putCopiedOrAwait(entry, publishToWriter);
1 ✔
438
  }
439

440
  /**
441
   * Associates the copied value with the copied key in the cache.
442
   *
443
   * @param entry the caller's entry alongside the copies to store
444
   * @param publishToWriter if the writer should be notified
445
   * @return the oldValue and if the new value was stored
446
   */
447
  @CanIgnoreReturnValue
448
  private PutResult<V> putCopiedOrAwait(CopiedEntry<K, V> entry, boolean publishToWriter) {
449
    K key = entry.key;
1 ✔
450
    V newValue = entry.copiedValue;
1 ✔
451
    var result = new PutResult<V>();
1 ✔
452
    BiFunction<K, @Nullable Expirable<V>, @Nullable Expirable<V>> remappingFunction =
1 ✔
453
        (K k, @Var @Nullable Expirable<V> expirable) -> {
454
      if (publishToWriter) {
1 ✔
455
        publishToCacheWriter(writer::write, () -> new EntryProxy<>(key, entry.value));
1 ✔
456
      }
457
      if ((expirable != null) && !expirable.isEternal()
1 ✔
458
          && expirable.hasExpired(currentTimeMillis())) {
1 ✔
459
        dispatcher.publishExpired(this, key, expirable.get());
1 ✔
460
        statistics.recordEvictions(1L);
1 ✔
461
        expirable = null;
1 ✔
462
      }
463
      @Var long expireTimeMillis = getWriteExpireTimeMillis((expirable == null));
1 ✔
464
      if ((expirable != null) && (expireTimeMillis == Long.MIN_VALUE)) {
1 ✔
465
        expireTimeMillis = expirable.getExpireTimeMillis();
1 ✔
466
      }
467
      if ((expireTimeMillis == 0) && (expirable == null)) {
1 ✔
468
        result.written = false;
1 ✔
469
        dispatcher.publishExpired(this, key, newValue);
1 ✔
470
        return null;
1 ✔
471
      } else if (expirable == null) {
1 ✔
472
        dispatcher.publishCreated(this, key, newValue);
1 ✔
473
      } else {
474
        result.oldValue = expirable.get();
1 ✔
475
        dispatcher.publishUpdated(this, key, expirable.get(), newValue);
1 ✔
476
      }
477
      result.written = true;
1 ✔
478
      return new Expirable<>(newValue, expireTimeMillis);
1 ✔
479
    };
480
    dispatcher.beginComputation();
1 ✔
481
    try {
482
      cache.asMap().compute(entry.copiedKey, remappingFunction);
1 ✔
483
    } finally {
484
      dispatcher.endComputation();
1 ✔
485
    }
486
    return result;
1 ✔
487
  }
488

489
  @Override
490
  @SuppressWarnings("CatchingUnchecked")
491
  public void putAll(Map<? extends K, ? extends V> map) {
492
    requireOperable();
1 ✔
493

494
    boolean statsEnabled = statistics.isEnabled();
1 ✔
495
    long start = statsEnabled ? ticker.read() : 0L;
1 ✔
496
    var copies = new ArrayList<CopiedEntry<K, V>>(map.size());
1 ✔
497
    for (var entry : map.entrySet()) {
1 ✔
498
      K key = requireNonNull(entry.getKey());
1 ✔
499
      V value = requireNonNull(entry.getValue());
1 ✔
500
      copies.add(new CopiedEntry<>(key, copyOf(key), value, copyOf(value)));
1 ✔
501
    }
1 ✔
502

503
    @Var CacheWriterException error = null;
1 ✔
504
    @Var Set<? extends K> failedKeys = Set.of();
1 ✔
505
    if (configuration.isWriteThrough() && !copies.isEmpty()) {
1 ✔
506
      var entries = new ArrayList<Cache.Entry<? extends K, ? extends V>>(copies.size());
1 ✔
507
      for (var copy : copies) {
1 ✔
508
        entries.add(new EntryProxy<>(copy.key, copy.value));
1 ✔
509
      }
1 ✔
510
      try {
511
        writer.writeAll(entries);
1 ✔
512
      } catch (CacheWriterException e) {
1 ✔
513
        failedKeys = entries.stream().map(Cache.Entry::getKey).collect(toSet());
1 ✔
514
        error = e;
1 ✔
515
      } catch (Exception e) {
1 ✔
516
        restoreInterrupt(e);
1 ✔
517
        failedKeys = entries.stream().map(Cache.Entry::getKey).collect(toSet());
1 ✔
518
        error = new CacheWriterException("Exception in CacheWriter", e);
1 ✔
519
      }
1 ✔
520
    }
521

522
    @Var int puts = 0;
1 ✔
523
    try {
524
      for (var copy : copies) {
1 ✔
525
        if (!failedKeys.contains(copy.key)) {
1 ✔
526
          var result = putCopiedOrAwait(copy, /* publishToWriter= */ false);
1 ✔
527
          if (result.written) {
1 ✔
528
            puts++;
1 ✔
529
          }
530
        }
531
      }
1 ✔
532
    } catch (Throwable t) {
1 ✔
533
      if (error != null) {
1 ✔
534
        t.addSuppressed(error);
1 ✔
535
      }
536
      awaitAndSuppressFailure(t);
1 ✔
537
      throw t;
1 ✔
538
    }
1 ✔
539
    var listenerFailure = awaitSynchronousFailure();
1 ✔
540

541
    if (statsEnabled && (puts > 0)) {
1 ✔
542
      statistics.recordPuts(puts);
1 ✔
543
      statistics.recordPutTime(ticker.read() - start);
1 ✔
544
    }
545
    var failure = suppress(error, listenerFailure);
1 ✔
546
    if (failure != null) {
1 ✔
547
      throw failure;
1 ✔
548
    }
549
  }
1 ✔
550

551
  @Override
552
  public boolean putIfAbsent(K key, V value) {
553
    requireOperable();
1 ✔
554
    requireNonNull(value);
1 ✔
555

556
    boolean statsEnabled = statistics.isEnabled();
1 ✔
557
    long start = statsEnabled ? ticker.read() : 0L;
1 ✔
558
    var result = putIfAbsentNoAwait(key, value, /* publishToWriter= */ true);
1 ✔
559
    var listenerFailure = awaitSynchronousFailure();
1 ✔
560
    if (statsEnabled) {
1 ✔
561
      if (result.oldValue != null) {
1 ✔
562
        statistics.recordHits(1L);
1 ✔
563
      } else {
564
        statistics.recordMisses(1L);
1 ✔
565
      }
566
      if (result.written) {
1 ✔
567
        statistics.recordPuts(1L);
1 ✔
568
        statistics.recordPutTime(ticker.read() - start);
1 ✔
569
      }
570
    }
571
    rethrowListenerFailure(listenerFailure);
1 ✔
572
    return result.written;
1 ✔
573
  }
574

575
  /**
576
   * Associates the specified value with the specified key in the cache if there is no existing
577
   * mapping.
578
   *
579
   * @param key key with which the specified value is to be associated
580
   * @param value value to be associated with the specified key
581
   * @param publishToWriter if the writer should be notified
582
   * @return the oldValue and if the new value was stored
583
   */
584
  @CanIgnoreReturnValue
585
  private PutResult<V> putIfAbsentNoAwait(K key, V value, boolean publishToWriter) {
586
    var result = new PutResult<V>();
1 ✔
587
    BiFunction<K, @Nullable Expirable<V>, @Nullable Expirable<V>> remappingFunction =
1 ✔
588
        (K k, @Nullable Expirable<V> expirable) -> {
589
      if ((expirable != null)
1 ✔
590
          && (expirable.isEternal() || !expirable.hasExpired(currentTimeMillis()))) {
1 ✔
591
        result.oldValue = expirable.get();
1 ✔
592
        return expirable;
1 ✔
593
      }
594

595
      V copy = copyOf(value);
1 ✔
596
      if (publishToWriter) {
1 ✔
597
        publishToCacheWriter(writer::write, () -> new EntryProxy<>(key, value));
1 ✔
598
      }
599
      if (expirable != null) {
1 ✔
600
        dispatcher.publishExpired(this, key, expirable.get());
1 ✔
601
        statistics.recordEvictions(1L);
1 ✔
602
      }
603

604
      long expireTimeMillis = getWriteExpireTimeMillis(/* created= */ true);
1 ✔
605
      if (expireTimeMillis == 0) {
1 ✔
606
        // A zero creation expiry means the entry is already expired and is not added
607
        dispatcher.publishExpired(this, key, copy);
1 ✔
608
        return null;
1 ✔
609
      } else {
610
        result.written = true;
1 ✔
611
        dispatcher.publishCreated(this, key, copy);
1 ✔
612
        return new Expirable<>(copy, expireTimeMillis);
1 ✔
613
      }
614
    };
615
    dispatcher.beginComputation();
1 ✔
616
    try {
617
      cache.asMap().compute(copyOf(key), remappingFunction);
1 ✔
618
    } finally {
619
      dispatcher.endComputation();
1 ✔
620
    }
621
    return result;
1 ✔
622
  }
623

624
  @Override
625
  public boolean remove(K key) {
626
    requireOperable();
1 ✔
627
    requireNonNull(key);
1 ✔
628

629
    boolean statsEnabled = statistics.isEnabled();
1 ✔
630
    long start = statsEnabled ? ticker.read() : 0L;
1 ✔
631
    V value = removeNoCopyOrAwait(key, /* publishToWriter= */ true);
1 ✔
632
    var listenerFailure = awaitSynchronousFailure();
1 ✔
633
    if (value != null) {
1 ✔
634
      statistics.recordRemovals(1L);
1 ✔
635
      if (statsEnabled) {
1 ✔
636
        statistics.recordRemoveTime(ticker.read() - start);
1 ✔
637
      }
638
    }
639
    rethrowListenerFailure(listenerFailure);
1 ✔
640
    return (value != null);
1 ✔
641
  }
642

643
  /**
644
   * Removes the mapping from the cache without store-by-value copying nor waiting for synchronous
645
   * listeners to complete.
646
   *
647
   * @param key key whose mapping is to be removed from the cache
648
   * @param publishToWriter if the writer should be notified
649
   * @return the old value
650
   */
651
  private @Nullable V removeNoCopyOrAwait(K key, boolean publishToWriter) {
652
    @SuppressWarnings("unchecked")
653
    var removed = (V[]) new Object[1];
1 ✔
654
    BiFunction<K, @Nullable Expirable<V>, @Nullable Expirable<V>> remappingFunction =
1 ✔
655
        (K k, @Nullable Expirable<V> expirable) -> {
656
      if (publishToWriter) {
1 ✔
657
        publishToCacheWriter(writer::delete, () -> key);
1 ✔
658
      }
659
      if (expirable != null) {
1 ✔
660
        if (!expirable.isEternal() && expirable.hasExpired(currentTimeMillis())) {
1 ✔
661
          dispatcher.publishExpired(this, key, expirable.get());
1 ✔
662
          statistics.recordEvictions(1L);
1 ✔
663
        } else {
664
          dispatcher.publishRemoved(this, key, expirable.get());
1 ✔
665
          removed[0] = expirable.get();
1 ✔
666
        }
667
      }
668
      return null;
1 ✔
669
    };
670
    dispatcher.beginComputation();
1 ✔
671
    try {
672
      cache.asMap().compute(key, remappingFunction);
1 ✔
673
    } finally {
674
      dispatcher.endComputation();
1 ✔
675
    }
676
    return removed[0];
1 ✔
677
  }
678

679
  @Override
680
  @CanIgnoreReturnValue
681
  public boolean remove(K key, V oldValue) {
682
    requireOperable();
1 ✔
683
    requireNonNull(key);
1 ✔
684
    requireNonNull(oldValue);
1 ✔
685

686
    boolean statsEnabled = statistics.isEnabled();
1 ✔
687
    long start = statsEnabled ? ticker.read() : 0L;
1 ✔
688
    boolean[] removed = { false };
1 ✔
689
    boolean[] found = { false };
1 ✔
690
    dispatcher.beginComputation();
1 ✔
691
    try {
692
      cache.asMap().computeIfPresent(key, (k, expirable) -> {
1 ✔
693
        long millis = expirable.isEternal()
1 ✔
694
            ? 0L
1 ✔
695
            : nanosToMillis((start == 0L) ? ticker.read() : start);
1 ✔
696
        if (expirable.hasExpired(millis)) {
1 ✔
697
          dispatcher.publishExpired(this, key, expirable.get());
1 ✔
698
          statistics.recordEvictions(1L);
1 ✔
699
          return null;
1 ✔
700
        }
701

702
        found[0] = true;
1 ✔
703
        if (oldValue.equals(expirable.get())) {
1 ✔
704
          publishToCacheWriter(writer::delete, () -> key);
1 ✔
705
          dispatcher.publishRemoved(this, key, expirable.get());
1 ✔
706
          removed[0] = true;
1 ✔
707
          return null;
1 ✔
708
        }
709
        setAccessExpireTime(expirable, getAccessExpireTime(), millis);
1 ✔
710
        return expirable;
1 ✔
711
      });
712
    } finally {
713
      dispatcher.endComputation();
1 ✔
714
    }
715
    var listenerFailure = awaitSynchronousFailure();
1 ✔
716
    if (statsEnabled) {
1 ✔
717
      if (removed[0]) {
1 ✔
718
        statistics.recordRemovals(1L);
1 ✔
719
        statistics.recordHits(1L);
1 ✔
720
        statistics.recordRemoveTime(ticker.read() - start);
1 ✔
721
      } else if (found[0]) {
1 ✔
722
        statistics.recordHits(1L);
1 ✔
723
      } else {
724
        statistics.recordMisses(1L);
1 ✔
725
      }
726
    }
727
    rethrowListenerFailure(listenerFailure);
1 ✔
728
    return removed[0];
1 ✔
729
  }
730

731
  @Override
732
  public @Nullable V getAndRemove(K key) {
733
    requireOperable();
1 ✔
734
    requireNonNull(key);
1 ✔
735

736
    boolean statsEnabled = statistics.isEnabled();
1 ✔
737
    long start = statsEnabled ? ticker.read() : 0L;
1 ✔
738
    V value = removeNoCopyOrAwait(key, /* publishToWriter= */ true);
1 ✔
739
    var listenerFailure = awaitSynchronousFailure();
1 ✔
740
    if (statsEnabled) {
1 ✔
741
      long duration = ticker.read() - start;
1 ✔
742
      if (value == null) {
1 ✔
743
        statistics.recordMisses(1L);
1 ✔
744
      } else {
745
        statistics.recordHits(1L);
1 ✔
746
        statistics.recordRemovals(1L);
1 ✔
747
        statistics.recordRemoveTime(duration);
1 ✔
748
      }
749
      statistics.recordGetTime(duration);
1 ✔
750
    }
751
    rethrowListenerFailure(listenerFailure);
1 ✔
752
    return (value == null) ? null : copyOf(value);
1 ✔
753
  }
754

755
  @Override
756
  public boolean replace(K key, V oldValue, V newValue) {
757
    requireOperable();
1 ✔
758
    requireNonNull(oldValue);
1 ✔
759
    requireNonNull(newValue);
1 ✔
760

761
    boolean statsEnabled = statistics.isEnabled();
1 ✔
762
    long start = statsEnabled ? ticker.read() : 0L;
1 ✔
763
    boolean[] replaced = { false };
1 ✔
764
    boolean[] found = { false };
1 ✔
765
    dispatcher.beginComputation();
1 ✔
766
    try {
767
      cache.asMap().computeIfPresent(key, (k, expirable) -> {
1 ✔
768
        long millis = expirable.isEternal()
1 ✔
769
            ? 0L
1 ✔
770
            : nanosToMillis((start == 0L) ? ticker.read() : start);
1 ✔
771
        if (expirable.hasExpired(millis)) {
1 ✔
772
          dispatcher.publishExpired(this, key, expirable.get());
1 ✔
773
          statistics.recordEvictions(1L);
1 ✔
774
          return null;
1 ✔
775
        }
776

777
        found[0] = true;
1 ✔
778
        Expirable<V> result;
779
        if (oldValue.equals(expirable.get())) {
1 ✔
780
          V copy = copyOf(newValue);
1 ✔
781
          publishToCacheWriter(writer::write, () -> new EntryProxy<>(key, newValue));
1 ✔
782
          dispatcher.publishUpdated(this, key, expirable.get(), copy);
1 ✔
783
          @Var long expireTimeMillis = getWriteExpireTimeMillis(/* created= */ false);
1 ✔
784
          if (expireTimeMillis == Long.MIN_VALUE) {
1 ✔
785
            expireTimeMillis = expirable.getExpireTimeMillis();
1 ✔
786
          }
787
          result = new Expirable<>(copy, expireTimeMillis);
1 ✔
788
          replaced[0] = true;
1 ✔
789
        } else {
1 ✔
790
          result = expirable;
1 ✔
791
          setAccessExpireTime(expirable, getAccessExpireTime(), millis);
1 ✔
792
        }
793
        return result;
1 ✔
794
      });
795
    } finally {
796
      dispatcher.endComputation();
1 ✔
797
    }
798
    var listenerFailure = awaitSynchronousFailure();
1 ✔
799

800
    if (statsEnabled) {
1 ✔
801
      statistics.recordMisses(found[0] ? 0L : 1L);
1 ✔
802
      statistics.recordHits(found[0] ? 1L : 0L);
1 ✔
803
      long duration = ticker.read() - start;
1 ✔
804
      if (replaced[0]) {
1 ✔
805
        statistics.recordPuts(1L);
1 ✔
806
        statistics.recordPutTime(duration);
1 ✔
807
      }
808
      statistics.recordGetTime(duration);
1 ✔
809
    }
810

811
    rethrowListenerFailure(listenerFailure);
1 ✔
812
    return replaced[0];
1 ✔
813
  }
814

815
  @Override
816
  public boolean replace(K key, V value) {
817
    requireOperable();
1 ✔
818

819
    boolean statsEnabled = statistics.isEnabled();
1 ✔
820
    long start = statsEnabled ? ticker.read() : 0L;
1 ✔
821
    var oldValue = replaceNoCopyOrAwait(key, value);
1 ✔
822
    var listenerFailure = awaitSynchronousFailure();
1 ✔
823
    if (oldValue == null) {
1 ✔
824
      statistics.recordMisses(1L);
1 ✔
825
      if (statsEnabled) {
1 ✔
826
        statistics.recordGetTime(ticker.read() - start);
1 ✔
827
      }
828
    } else if (statsEnabled) {
1 ✔
829
      statistics.recordHits(1L);
1 ✔
830
      statistics.recordPuts(1L);
1 ✔
831
      long duration = ticker.read() - start;
1 ✔
832
      statistics.recordGetTime(duration);
1 ✔
833
      statistics.recordPutTime(duration);
1 ✔
834
    }
835
    rethrowListenerFailure(listenerFailure);
1 ✔
836
    return (oldValue != null);
1 ✔
837
  }
838

839
  @Override
840
  public @Nullable V getAndReplace(K key, V value) {
841
    requireOperable();
1 ✔
842

843
    boolean statsEnabled = statistics.isEnabled();
1 ✔
844
    long start = statsEnabled ? ticker.read() : 0L;
1 ✔
845
    V oldValue = replaceNoCopyOrAwait(key, value);
1 ✔
846
    var listenerFailure = awaitSynchronousFailure();
1 ✔
847
    if (statsEnabled) {
1 ✔
848
      long duration = ticker.read() - start;
1 ✔
849
      if (oldValue == null) {
1 ✔
850
        statistics.recordMisses(1L);
1 ✔
851
      } else {
852
        statistics.recordHits(1L);
1 ✔
853
        statistics.recordPuts(1L);
1 ✔
854
        statistics.recordPutTime(duration);
1 ✔
855
      }
856
      statistics.recordGetTime(duration);
1 ✔
857
    }
858
    rethrowListenerFailure(listenerFailure);
1 ✔
859
    return (oldValue == null) ? null : copyOf(oldValue);
1 ✔
860
  }
861

862
  /**
863
   * Replaces the entry for the specified key only if it is currently mapped to some value. The
864
   * entry is not store-by-value copied nor does the method wait for synchronous listeners to
865
   * complete.
866
   *
867
   * @param key key with which the specified value is associated
868
   * @param value value to be associated with the specified key
869
   * @return the old value
870
   */
871
  private @Nullable V replaceNoCopyOrAwait(K key, V value) {
872
    requireNonNull(value);
1 ✔
873
    @SuppressWarnings("unchecked")
874
    var replaced = (V[]) new Object[1];
1 ✔
875
    dispatcher.beginComputation();
1 ✔
876
    try {
877
      cache.asMap().computeIfPresent(key, (k, expirable) -> {
1 ✔
878
        if (!expirable.isEternal() && expirable.hasExpired(currentTimeMillis())) {
1 ✔
879
          dispatcher.publishExpired(this, key, expirable.get());
1 ✔
880
          statistics.recordEvictions(1L);
1 ✔
881
          return null;
1 ✔
882
        }
883

884
        V copy = copyOf(value);
1 ✔
885
        publishToCacheWriter(writer::write, () -> new EntryProxy<>(key, value));
1 ✔
886
        @Var long expireTimeMillis = getWriteExpireTimeMillis(/* created= */ false);
1 ✔
887
        if (expireTimeMillis == Long.MIN_VALUE) {
1 ✔
888
          expireTimeMillis = expirable.getExpireTimeMillis();
1 ✔
889
        }
890
        dispatcher.publishUpdated(this, key, expirable.get(), copy);
1 ✔
891
        replaced[0] = expirable.get();
1 ✔
892
        return new Expirable<>(copy, expireTimeMillis);
1 ✔
893
      });
894
    } finally {
895
      dispatcher.endComputation();
1 ✔
896
    }
897
    return replaced[0];
1 ✔
898
  }
899

900
  @Override
901
  @SuppressWarnings("CatchingUnchecked")
902
  public void removeAll(Set<? extends K> keys) {
903
    requireOperable();
1 ✔
904
    var keysToRemove = new LinkedHashSet<>(keys);
1 ✔
905
    keysToRemove.forEach(Objects::requireNonNull);
1 ✔
906

907
    @Var CacheWriterException error = null;
1 ✔
908
    @Var Set<? extends K> failedKeys = Set.of();
1 ✔
909
    boolean statsEnabled = statistics.isEnabled();
1 ✔
910
    long start = statsEnabled ? ticker.read() : 0L;
1 ✔
911
    if (configuration.isWriteThrough() && !keysToRemove.isEmpty()) {
1 ✔
912
      var keysToWrite = new LinkedHashSet<>(keysToRemove);
1 ✔
913
      try {
914
        writer.deleteAll(keysToWrite);
1 ✔
915
      } catch (CacheWriterException e) {
1 ✔
916
        error = e;
1 ✔
917
        failedKeys = keysToWrite;
1 ✔
918
      } catch (Exception e) {
1 ✔
919
        restoreInterrupt(e);
1 ✔
920
        error = new CacheWriterException("Exception in CacheWriter", e);
1 ✔
921
        failedKeys = keysToWrite;
1 ✔
922
      }
1 ✔
923
    }
924

925
    @Var int removed = 0;
1 ✔
926
    try {
927
      for (var key : keysToRemove) {
1 ✔
928
        if (!failedKeys.contains(key)
1 ✔
929
            && (removeNoCopyOrAwait(key, /* publishToWriter= */ false) != null)) {
1 ✔
930
          removed++;
1 ✔
931
        }
932
      }
1 ✔
933
    } catch (Throwable t) {
1 ✔
934
      if (error != null) {
1 !
UNCOV
935
        t.addSuppressed(error);
×
936
      }
937
      awaitAndSuppressFailure(t);
1 ✔
938
      throw t;
1 ✔
939
    }
1 ✔
940
    var listenerFailure = awaitSynchronousFailure();
1 ✔
941

942
    if (statsEnabled && (removed > 0)) {
1 ✔
943
      statistics.recordRemovals(removed);
1 ✔
944
      statistics.recordRemoveTime(ticker.read() - start);
1 ✔
945
    }
946
    var failure = suppress(error, listenerFailure);
1 ✔
947
    if (failure != null) {
1 ✔
948
      throw failure;
1 ✔
949
    }
950
  }
1 ✔
951

952
  @Override
953
  public void removeAll() {
954
    removeAll(cache.asMap().keySet());
1 ✔
955
  }
1 ✔
956

957
  @Override
958
  public void clear() {
959
    requireOperable();
1 ✔
960
    dispatcher.beginComputation();
1 ✔
961
    try {
962
      cache.invalidateAll();
1 ✔
963
    } finally {
964
      dispatcher.endComputation();
1 ✔
965
    }
966
  }
1 ✔
967

968
  @Override
969
  public <C extends Configuration<K, V>> C getConfiguration(Class<C> clazz) {
970
    if (clazz.isInstance(configuration)) {
1 ✔
971
      synchronized (configuration) {
1 ✔
972
        return clazz.cast(configuration.immutableCopy());
1 ✔
973
      }
974
    }
975
    throw new IllegalArgumentException("The configuration class " + clazz
1 ✔
976
        + " is not supported by this implementation");
977
  }
978

979
  @Override
980
  public <T extends @Nullable Object> T invoke(K key,
981
      EntryProcessor<K, V, T> entryProcessor, Object @Nullable ... arguments) {
982
    requireNonNull(entryProcessor);
1 ✔
983
    requireOperable();
1 ✔
984

985
    var result = new Object[1];
1 ✔
986
    var failure = new Throwable[1];
1 ✔
987
    BiFunction<K, @Nullable Expirable<V>, @Nullable Expirable<V>> remappingFunction =
1 ✔
988
        (k, expirable) -> {
989
      // Publish a lazily-expired prior's expiration before the processor observes it as absent,
990
      // so listeners see a linearizable sequence and the expiration is committed exactly once
991
      boolean expired;
992
      Expirable<V> prior;
993
      if ((expirable != null) && !expirable.isEternal()
1 ✔
994
          && expirable.hasExpired(currentTimeMillis())) {
1 ✔
995
        dispatcher.publishExpired(this, key, expirable.get());
1 ✔
996
        statistics.recordEvictions(1L);
1 ✔
997
        expired = true;
1 ✔
998
        prior = null;
1 ✔
999
      } else {
1000
        prior = expirable;
1 ✔
1001
        expired = false;
1 ✔
1002
      }
1003

1004
      V value;
1005
      if (prior == null) {
1 ✔
1006
        statistics.recordMisses(1L);
1 ✔
1007
        value = null;
1 ✔
1008
      } else {
1009
        value = copyOf(prior.get());
1 ✔
1010
        statistics.recordHits(1L);
1 ✔
1011
      }
1012
      var entry = new EntryProcessorEntry<>(key, value,
1 ✔
1013
          configuration.isReadThrough() ? cacheLoader : Optional.empty());
1 ✔
1014
      try {
1015
        result[0] = entryProcessor.process(entry, arguments);
1 ✔
1016
        return postProcess(prior, entry);
1 ✔
1017
      } catch (Throwable e) {
1 ✔
1018
        if (!expired) {
1 ✔
1019
          throw processorFailure(e);
1 ✔
1020
        }
1021
        failure[0] = e;
1 ✔
1022
        return null;
1 ✔
1023
      }
1024
    };
1025
    try {
1026
      dispatcher.beginComputation();
1 ✔
1027
      try {
1028
        cache.asMap().compute(copyOf(key), remappingFunction);
1 ✔
1029
      } finally {
1030
        dispatcher.endComputation();
1 ✔
1031
      }
1032
    } catch (Throwable t) {
1 ✔
1033
      dispatcher.ignoreSynchronous();
1 ✔
1034
      throw t;
1 ✔
1035
    }
1 ✔
1036
    var listenerFailure = awaitSynchronousFailure();
1 ✔
1037
    if (failure[0] != null) {
1 ✔
1038
      var error = (failure[0] instanceof Error) ? failure[0] : processorFailure(failure[0]);
1 ✔
1039
      if (listenerFailure != null) {
1 ✔
1040
        error.addSuppressed(listenerFailure);
1 ✔
1041
      }
1042
      throw processorFailure(error);
1 ✔
1043
    }
1044
    rethrowListenerFailure(listenerFailure);
1 ✔
1045

1046
    @SuppressWarnings("unchecked")
1047
    var castedResult = (T) result[0];
1 ✔
1048
    return castedResult;
1 ✔
1049
  }
1050

1051
  /**
1052
   * Waits for the synchronous listeners and returns their failure instead of throwing it. An
1053
   * {@code Error} is not captured and propagates.
1054
   */
1055
  protected final @Nullable CacheEntryListenerException awaitSynchronousFailure() {
1056
    try {
1057
      dispatcher.awaitSynchronous();
1 ✔
1058
      return null;
1 ✔
1059
    } catch (CacheEntryListenerException e) {
1 ✔
1060
      return e;
1 ✔
1061
    }
1062
  }
1063

1064
  /** Waits for the synchronous listeners and suppresses any failures onto the existing error. */
1065
  @SuppressWarnings("CheckReturnValue")
1066
  protected final void awaitAndSuppressFailure(Throwable error) {
1067
    suppress(error, awaitSynchronousFailure());
1 ✔
1068
  }
1 ✔
1069

1070
  /** Rethrows a captured synchronous listener failure after the operation is accounted for. */
1071
  protected static void rethrowListenerFailure(@Nullable CacheEntryListenerException failure) {
1072
    if (failure != null) {
1 ✔
1073
      throw failure;
1 ✔
1074
    }
1075
  }
1 ✔
1076

1077
  /** Returns the exception to rethrow on an entry processor failure. */
1078
  private static RuntimeException processorFailure(Throwable e) {
1079
    if (e instanceof Error) {
1 ✔
1080
      throw (Error) e;
1 ✔
1081
    } else if (e instanceof EntryProcessorException) {
1 ✔
1082
      return (EntryProcessorException) e;
1 ✔
1083
    }
1084
    return new EntryProcessorException(e);
1 ✔
1085
  }
1086

1087
  /** Restores the thread's interrupt status if the failure is an interruption. */
1088
  private static void restoreInterrupt(Throwable failure) {
1089
    if (failure instanceof InterruptedException) {
1 ✔
1090
      Thread.currentThread().interrupt();
1 ✔
1091
    }
1092
  }
1 ✔
1093

1094
  /**
1095
   * Returns the updated expirable value after performing the post-processing actions. A null
1096
   * {@code expirable} means the entry was absent and READ/UPDATED (which require a live prior)
1097
   * never observe one.
1098
   */
1099
  @Nullable Expirable<V> postProcess(
1100
      @Nullable Expirable<V> expirable, EntryProcessorEntry<K, V> entry) {
1101
    switch (entry.getAction()) {
1 ✔
1102
      case NONE:
1103
        return expirable;
1 ✔
1104
      case READ: {
1105
        setAccessExpireTime(requireNonNull(expirable), getAccessExpireTime(), 0L);
1 ✔
1106
        return expirable;
1 ✔
1107
      }
1108
      case CREATED:
1109
      case LOADED: {
1110
        V value = requireNonNull(entry.getValue());
1 ✔
1111
        V copy = copyOf(value);
1 ✔
1112
        if (entry.getAction() == Action.CREATED) {
1 ✔
1113
          publishToCacheWriter(writer::write, () -> new EntryProxy<>(entry.getKey(), value));
1 ✔
1114
        }
1115
        long expireTimeMillis = getWriteExpireTimeMillis(/* created= */ true);
1 ✔
1116
        if (expireTimeMillis == 0) {
1 ✔
1117
          // A zero creation expiry means the entry is already expired and is not added, so the
1118
          // create is neither counted nor published
1119
          dispatcher.publishExpired(this, entry.getKey(), copy);
1 ✔
1120
          return null;
1 ✔
1121
        }
1122
        if (entry.getAction() == Action.CREATED) {
1 ✔
1123
          statistics.recordPuts(1L);
1 ✔
1124
        }
1125
        dispatcher.publishCreated(this, entry.getKey(), copy);
1 ✔
1126
        return new Expirable<>(copy, expireTimeMillis);
1 ✔
1127
      }
1128
      case UPDATED: {
1129
        requireNonNull(expirable, "Expected a previous value but was null");
1 ✔
1130
        V value = requireNonNull(entry.getValue(), "Expected a new value but was null");
1 ✔
1131
        V copy = copyOf(value);
1 ✔
1132
        publishToCacheWriter(writer::write, () -> new EntryProxy<>(entry.getKey(), value));
1 ✔
1133
        statistics.recordPuts(1L);
1 ✔
1134
        dispatcher.publishUpdated(this, entry.getKey(), expirable.get(), copy);
1 ✔
1135
        @Var long expireTimeMillis = getWriteExpireTimeMillis(/* created= */ false);
1 ✔
1136
        if (expireTimeMillis == Long.MIN_VALUE) {
1 ✔
1137
          expireTimeMillis = expirable.getExpireTimeMillis();
1 ✔
1138
        }
1139
        return new Expirable<>(copy, expireTimeMillis);
1 ✔
1140
      }
1141
      case DELETED:
1142
        publishToCacheWriter(writer::delete, entry::getKey);
1 ✔
1143
        if (expirable != null) {
1 ✔
1144
          statistics.recordRemovals(1L);
1 ✔
1145
          dispatcher.publishRemoved(this, entry.getKey(), expirable.get());
1 ✔
1146
        }
1147
        return null;
1 ✔
1148
    }
1149
    throw new IllegalStateException("Unknown state: " + entry.getAction());
1 ✔
1150
  }
1151

1152
  @Override
1153
  public <T> Map<K, EntryProcessorResult<T>> invokeAll(Set<? extends K> keys,
1154
      EntryProcessor<K, V, T> entryProcessor, Object @Nullable ... arguments) {
1155
    requireOperable();
1 ✔
1156
    requireNonNull(keys);
1 ✔
1157
    requireNonNull(entryProcessor);
1 ✔
1158
    keys.forEach(Objects::requireNonNull);
1 ✔
1159

1160
    var results = new HashMap<K, EntryProcessorResult<T>>(keys.size(), 1.0f);
1 ✔
1161
    for (K key : keys) {
1 ✔
1162
      try {
1163
        T result = invoke(key, entryProcessor, arguments);
1 ✔
1164
        if (result != null) {
1 ✔
1165
          results.put(key, () -> result);
1 ✔
1166
        }
1167
      } catch (EntryProcessorException e) {
1 ✔
1168
        results.put(key, () -> { throw e; });
1 ✔
1169
      } catch (RuntimeException e) {
1 ✔
1170
        results.put(key, () -> { throw new EntryProcessorException(e); });
1 ✔
1171
      }
1 ✔
1172
    }
1 ✔
1173
    return results;
1 ✔
1174
  }
1175

1176
  @Override
1177
  public String getName() {
1178
    return name;
1 ✔
1179
  }
1180

1181
  @Override
1182
  public CacheManager getCacheManager() {
1183
    return cacheManager;
1 ✔
1184
  }
1185

1186
  @Override
1187
  public boolean isClosed() {
1188
    return closed;
1 ✔
1189
  }
1190

1191
  @Override
1192
  public void close() {
1193
    if (isClosed()) {
1 ✔
1194
      return;
1 ✔
1195
    }
1196
    @Var Throwable thrown = null;
1 ✔
1197
    synchronized (configuration) {
1 ✔
1198
      if (isClosed()) {
1 ✔
1199
        return;
1 ✔
1200
      }
1201
      thrown = tryClose((AutoCloseable) () -> enableManagement(false), thrown);
1 ✔
1202
      thrown = tryClose((AutoCloseable) () -> enableStatistics(false), thrown);
1 ✔
1203

1204
      closed = true;
1 ✔
1205
      try {
1206
        cacheManager.destroyCache(name, this);
1 ✔
1207
      } catch (IllegalStateException ignored) { /* manager already closed */ }
1 ✔
1208
    }
1 ✔
1209

1210
    thrown = shutdownExecutor(thrown);
1 ✔
1211
    thrown = tryClose(expiry, thrown);
1 ✔
1212
    thrown = tryClose(writer, thrown);
1 ✔
1213
    thrown = tryClose(cacheLoader.orElse(null), thrown);
1 ✔
1214
    for (Registration<K, V> registration : dispatcher.registrations()) {
1 ✔
1215
      thrown = tryClose(registration.getCacheEntryListener(), thrown);
1 ✔
1216
    }
1 ✔
1217
    if (thrown != null) {
1 ✔
1218
      logger.log(Level.WARNING, "Failure when closing cache resources", thrown);
1 ✔
1219
    }
1220
    dispatcher.beginComputation();
1 ✔
1221
    try {
1222
      cache.invalidateAll();
1 ✔
1223
    } finally {
1224
      dispatcher.endComputation();
1 ✔
1225
    }
1226
  }
1 ✔
1227

1228
  @SuppressWarnings("FutureReturnValueIgnored")
1229
  private @Nullable Throwable shutdownExecutor(@Var @Nullable Throwable thrown) {
1230
    if (executor instanceof ExecutorService) {
1 ✔
1231
      @SuppressWarnings("PMD.CloseResource")
1232
      var es = (ExecutorService) executor;
1 ✔
1233
      thrown = tryClose((AutoCloseable) es::shutdown, thrown);
1 ✔
1234
    }
1235

1236
    try {
1237
      CompletableFuture
1 ✔
1238
          .allOf(inFlight.toArray(CompletableFuture[]::new))
1 ✔
1239
          .get(10, TimeUnit.SECONDS);
1 ✔
1240
    } catch (ExecutionException | TimeoutException e) {
1 ✔
1241
      thrown = suppress(thrown, e);
1 ✔
1242
    } catch (InterruptedException e) {
1 ✔
1243
      Thread.currentThread().interrupt();
1 ✔
1244
      thrown = suppress(thrown, e);
1 ✔
1245
    }
1 ✔
1246
    inFlight.clear();
1 ✔
1247

1248
    if (!(executor instanceof ExecutorService)) {
1 ✔
1249
      thrown = tryClose(executor, thrown);
1 ✔
1250
    }
1251
    return thrown;
1 ✔
1252
  }
1253

1254
  /**
1255
   * Attempts to close the resource. If an error occurs and an outermost exception is set, then adds
1256
   * the error to the suppression list.
1257
   *
1258
   * @param o the resource to close if Closeable
1259
   * @param outer the outermost error, or null if unset
1260
   * @return the outermost error, or null if unset and successful
1261
   */
1262
  private static @Nullable Throwable tryClose(@Nullable Object o, @Nullable Throwable outer) {
1263
    if (o instanceof AutoCloseable) {
1 ✔
1264
      try {
1265
        ((AutoCloseable) o).close();
1 ✔
1266
      } catch (Throwable t) {
1 ✔
1267
        return suppress(outer, t);
1 ✔
1268
      }
1 ✔
1269
    }
1270
    return outer;
1 ✔
1271
  }
1272

1273
  /** Returns the outermost error, retaining the error as suppressed if one is already set. */
1274
  private static <T extends Throwable> @Nullable T suppress(@Nullable T outer, @Nullable T t) {
1275
    if (t == null) {
1 ✔
1276
      return outer;
1 ✔
1277
    }
1278
    if (outer == null) {
1 ✔
1279
      return t;
1 ✔
1280
    }
1281
    if (outer != t) {
1 ✔
1282
      outer.addSuppressed(t);
1 ✔
1283
    }
1284
    return outer;
1 ✔
1285
  }
1286

1287
  @Override
1288
  public <T> T unwrap(Class<T> clazz) {
1289
    if (clazz.isInstance(cache)) {
1 ✔
1290
      return clazz.cast(cache);
1 ✔
1291
    } else if (clazz.isInstance(this)) {
1 ✔
1292
      return clazz.cast(this);
1 ✔
1293
    }
1294
    throw new IllegalArgumentException("Unwrapping to " + clazz
1 ✔
1295
        + " is not supported by this implementation");
1296
  }
1297

1298
  @Override
1299
  public void registerCacheEntryListener(
1300
      CacheEntryListenerConfiguration<K, V> cacheEntryListenerConfiguration) {
1301
    synchronized (configuration) {
1 ✔
1302
      requireOperable();
1 ✔
1303
      configuration.addCacheEntryListenerConfiguration(cacheEntryListenerConfiguration);
1 ✔
1304
      try {
1305
        dispatcher.register(cacheEntryListenerConfiguration);
1 ✔
1306
      } catch (Throwable t) {
1 ✔
1307
        configuration.removeCacheEntryListenerConfiguration(cacheEntryListenerConfiguration);
1 ✔
1308
        throw t;
1 ✔
1309
      }
1 ✔
1310
    }
1 ✔
1311
  }
1 ✔
1312

1313
  @Override
1314
  public void deregisterCacheEntryListener(
1315
      CacheEntryListenerConfiguration<K, V> cacheEntryListenerConfiguration) {
1316
    synchronized (configuration) {
1 ✔
1317
      requireOperable();
1 ✔
1318
      configuration.removeCacheEntryListenerConfiguration(cacheEntryListenerConfiguration);
1 ✔
1319
      dispatcher.deregister(cacheEntryListenerConfiguration);
1 ✔
1320
    }
1 ✔
1321
  }
1 ✔
1322

1323
  @Override
1324
  public Iterator<Cache.Entry<K, V>> iterator() {
1325
    requireOperable();
1 ✔
1326
    return new EntryIterator();
1 ✔
1327
  }
1328

1329
  /** Enables or disables the configuration management JMX bean. */
1330
  void enableManagement(boolean enabled) {
1331
    synchronized (configuration) {
1 ✔
1332
      requireNotClosed();
1 ✔
1333
      if (enabled) {
1 ✔
1334
        JmxRegistration.registerMxBean(this, cacheMxBean, MBeanType.CONFIGURATION);
1 ✔
1335
      } else if (configuration.isManagementEnabled()) {
1 ✔
1336
        JmxRegistration.unregisterMxBean(this, MBeanType.CONFIGURATION);
1 ✔
1337
      }
1338
      configuration.setManagementEnabled(enabled);
1 ✔
1339
    }
1 ✔
1340
  }
1 ✔
1341

1342
  /** Enables or disables the statistics JMX bean. */
1343
  void enableStatistics(boolean enabled) {
1344
    synchronized (configuration) {
1 ✔
1345
      requireNotClosed();
1 ✔
1346
      if (enabled) {
1 ✔
1347
        JmxRegistration.registerMxBean(this, statistics, MBeanType.STATISTICS);
1 ✔
1348
      } else if (configuration.isStatisticsEnabled()) {
1 ✔
1349
        JmxRegistration.unregisterMxBean(this, MBeanType.STATISTICS);
1 ✔
1350
      }
1351
      statistics.enable(enabled);
1 ✔
1352
      configuration.setStatisticsEnabled(enabled);
1 ✔
1353
    }
1 ✔
1354
  }
1 ✔
1355

1356
  /** Performs the action with the cache writer if write-through is enabled. */
1357
  @SuppressWarnings("CatchingUnchecked")
1358
  private <T> void publishToCacheWriter(Consumer<T> action, Supplier<T> data) {
1359
    if (!configuration.isWriteThrough()) {
1 ✔
1360
      return;
1 ✔
1361
    }
1362
    try {
1363
      action.accept(data.get());
1 ✔
1364
    } catch (CacheWriterException e) {
1 ✔
1365
      throw e;
1 ✔
1366
    } catch (Exception e) {
1 ✔
1367
      restoreInterrupt(e);
1 ✔
1368
      throw new CacheWriterException("Exception in CacheWriter", e);
1 ✔
1369
    }
1 ✔
1370
  }
1 ✔
1371

1372
  /** Checks that the cache is not closed. */
1373
  protected final void requireNotClosed() {
1374
    if (isClosed()) {
1 ✔
1375
      throw new IllegalStateException();
1 ✔
1376
    }
1377
  }
1 ✔
1378

1379
  /** Checks that the operation may begin. */
1380
  protected final void requireOperable() {
1381
    requireNotClosed();
1 ✔
1382
    if (dispatcher.inCallback()) {
1 ✔
1383
      throw new IllegalStateException("Recursive cache operation");
1 ✔
1384
    }
1385
  }
1 ✔
1386

1387
  /**
1388
   * Returns a copy of the value if value-based caching is enabled.
1389
   *
1390
   * @param object the object to be copied
1391
   * @param <T> the type of object being copied
1392
   * @return a copy of the object if storing by value or the same instance if by reference
1393
   */
1394
  @SuppressWarnings("CatchingUnchecked")
1395
  protected final <T> T copyOf(T object) {
1396
    try {
1397
      return requireNonNull(
1 ✔
1398
          copier.copy(requireNonNull(object), requireNonNull(cacheManager.getClassLoader())));
1 ✔
1399
    } catch (NullPointerException | IllegalStateException | ClassCastException | CacheException e) {
1 ✔
1400
      throw e;
1 ✔
1401
    } catch (Exception e) {
1 ✔
1402
      restoreInterrupt(e);
1 ✔
1403
      throw new CacheException(e);
1 ✔
1404
    }
1405
  }
1406

1407
  /**
1408
   * Returns a deep copy of the map if value-based caching is enabled.
1409
   *
1410
   * @param map the mapping of keys to expirable values
1411
   * @return a deep or shallow copy of the mappings depending on the store by value setting
1412
   */
1413
  @SuppressWarnings("CollectorMutability")
1414
  protected final Map<K, V> copyMap(Map<K, Expirable<V>> map) {
1415
    return map.entrySet().stream().collect(toMap(
1 ✔
1416
        entry -> copyOf(entry.getKey()),
1 ✔
1417
        entry -> copyOf(entry.getValue().get())));
1 ✔
1418
  }
1419

1420
  /** Returns the current time in milliseconds. */
1421
  protected final long currentTimeMillis() {
1422
    return nanosToMillis(ticker.read());
1 ✔
1423
  }
1424

1425
  /** Returns the nanosecond time in milliseconds. */
1426
  protected static long nanosToMillis(long nanos) {
1427
    return TimeUnit.NANOSECONDS.toMillis(nanos);
1 ✔
1428
  }
1429

1430
  /** Returns the duration to expire an accessed entry after, or {@code null} if unchanged. */
1431
  @SuppressWarnings("CatchingUnchecked")
1432
  protected final @Nullable Duration getAccessExpireTime() {
1433
    try {
1434
      return expiry.getExpiryForAccess();
1 ✔
1435
    } catch (Exception e) {
1 ✔
1436
      restoreInterrupt(e);
1 ✔
1437
      logger.log(Level.WARNING, "Failed to get the policy's expiration time", e);
1 ✔
1438
      return null;
1 ✔
1439
    }
1440
  }
1441

1442
  /**
1443
   * Sets the JCache access expiration time.
1444
   *
1445
   * @param expirable the entry that was operated on
1446
   * @param duration the access duration, or null if unchanged
1447
   * @param currentTimeMillis the current time, or zero if not read yet
1448
   */
1449
  protected final void setAccessExpireTime(Expirable<?> expirable,
1450
      @Nullable Duration duration, @Var long currentTimeMillis) {
1451
    if (duration == null) {
1 ✔
1452
      return;
1 ✔
1453
    } else if (duration.isZero()) {
1 ✔
1454
      expirable.setExpireTimeMillis(0L);
1 ✔
1455
    } else if (duration.isEternal()) {
1 ✔
1456
      expirable.setExpireTimeMillis(Long.MAX_VALUE);
1 ✔
1457
    } else {
1458
      if (currentTimeMillis == 0L) {
1 ✔
1459
        currentTimeMillis = currentTimeMillis();
1 ✔
1460
      }
1461
      @Var long expireTimeMillis = duration.getAdjustedTime(currentTimeMillis);
1 ✔
1462
      expireTimeMillis = ((expireTimeMillis == 0L) || (expireTimeMillis == Long.MAX_VALUE))
1 ✔
1463
          ? (expireTimeMillis - 1)
1 ✔
1464
          : expireTimeMillis;
1 ✔
1465
      expirable.setExpireTimeMillis(expireTimeMillis);
1 ✔
1466
    }
1467
  }
1 ✔
1468

1469
  /**
1470
   * Sets the native access expiration time.
1471
   *
1472
   * @param key the entry's key
1473
   * @param duration the access duration, or null if unchanged
1474
   */
1475
  protected final void setVariableExpiration(K key, @Nullable Duration duration) {
1476
    if (duration == null) {
1 ✔
1477
      return;
1 ✔
1478
    }
1479
    cache.policy().expireVariably().ifPresent(policy -> {
1 ✔
1480
      if (duration.isZero()) {
1 ✔
1481
        policy.setExpiresAfter(key, 0L, TimeUnit.NANOSECONDS);
1 ✔
1482
      } else if (duration.isEternal()) {
1 ✔
1483
        policy.setExpiresAfter(key, Long.MAX_VALUE, TimeUnit.NANOSECONDS);
1 ✔
1484
      } else {
1485
        policy.setExpiresAfter(key, duration.getDurationAmount(), duration.getTimeUnit());
1 ✔
1486
      }
1487
    });
1 ✔
1488
  }
1 ✔
1489

1490
  /**
1491
   * Returns the time when the entry will expire.
1492
   *
1493
   * @param created if the write operation is an insert or an update
1494
   * @return the time when the entry will expire, zero if it should expire immediately,
1495
   *         Long.MIN_VALUE if it should not be changed, or Long.MAX_VALUE if eternal
1496
   */
1497
  @SuppressWarnings("CatchingUnchecked")
1498
  protected final long getWriteExpireTimeMillis(boolean created) {
1499
    try {
1500
      Duration duration = created ? expiry.getExpiryForCreation() : expiry.getExpiryForUpdate();
1 ✔
1501
      if (duration == null) {
1 ✔
1502
        return created ? Long.MAX_VALUE : Long.MIN_VALUE;
1 ✔
1503
      } else if (duration.isZero()) {
1 ✔
1504
        return 0L;
1 ✔
1505
      } else if (duration.isEternal()) {
1 ✔
1506
        return Long.MAX_VALUE;
1 ✔
1507
      }
1508
      long expireTimeMillis = duration.getAdjustedTime(currentTimeMillis());
1 ✔
1509
      return ((expireTimeMillis == 0L) || (expireTimeMillis == Long.MAX_VALUE))
1 ✔
1510
          ? (expireTimeMillis - 1)
1 ✔
1511
          : expireTimeMillis;
1 ✔
1512
    } catch (Exception e) {
1 ✔
1513
      restoreInterrupt(e);
1 ✔
1514
      logger.log(Level.WARNING, "Failed to get the policy's expiration time", e);
1 ✔
1515
      return created ? Long.MAX_VALUE : Long.MIN_VALUE;
1 ✔
1516
    }
1517
  }
1518

1519
  /** An iterator to safely expose the cache entries. */
1520
  final class EntryIterator implements Iterator<Cache.Entry<K, V>> {
1 ✔
1521
    // NullAway does not yet understand the @NonNull annotation in the return type of asMap.
1522
    @SuppressWarnings("NullAway")
1 ✔
1523
    final Iterator<Map.Entry<K, Expirable<V>>> delegate = cache.asMap().entrySet().iterator();
1 ✔
1524

1525
    Map.@Nullable Entry<K, Expirable<V>> current;
1526
    Map.@Nullable Entry<K, Expirable<V>> cursor;
1527

1528
    @Override
1529
    public boolean hasNext() {
1530
      while ((cursor == null) && delegate.hasNext()) {
1 ✔
1531
        Map.Entry<K, Expirable<V>> entry = delegate.next();
1 ✔
1532
        long millis = entry.getValue().isEternal() ? 0L : currentTimeMillis();
1 ✔
1533
        if (!entry.getValue().hasExpired(millis)) {
1 ✔
1534
          var duration = getAccessExpireTime();
1 ✔
1535
          setAccessExpireTime(entry.getValue(), duration, millis);
1 ✔
1536
          setVariableExpiration(entry.getKey(), duration);
1 ✔
1537
          cursor = entry;
1 ✔
1538
        }
1539
      }
1 ✔
1540
      return (cursor != null);
1 ✔
1541
    }
1542

1543
    @Override
1544
    public Cache.Entry<K, V> next() {
1545
      if (!hasNext()) {
1 ✔
1546
        throw new NoSuchElementException();
1 ✔
1547
      }
1548
      var entry = requireNonNull(cursor);
1 ✔
1549
      cursor = null;
1 ✔
1550
      var copy = new EntryProxy<>(copyOf(entry.getKey()), copyOf(entry.getValue().get()));
1 ✔
1551
      statistics.recordHits(1L);
1 ✔
1552
      current = entry;
1 ✔
1553
      return copy;
1 ✔
1554
    }
1555

1556
    @Override
1557
    @SuppressWarnings("CheckReturnValue")
1558
    public void remove() {
1559
      if (current == null) {
1 ✔
1560
        throw new IllegalStateException();
1 ✔
1561
      }
1562
      CacheProxy.this.remove(current.getKey());
1 ✔
1563
      current = null;
1 ✔
1564
    }
1 ✔
1565
  }
1566

1567
  /** An entry paired with the store-by-value copies to be stored in its place. */
1568
  private static final class CopiedEntry<K, V> {
1569
    final K key;
1570
    final V value;
1571
    final K copiedKey;
1572
    final V copiedValue;
1573

1574
    CopiedEntry(K key, K copiedKey, V value, V copiedValue) {
1 ✔
1575
      this.copiedValue = copiedValue;
1 ✔
1576
      this.copiedKey = copiedKey;
1 ✔
1577
      this.value = value;
1 ✔
1578
      this.key = key;
1 ✔
1579
    }
1 ✔
1580
  }
1581

1582
  protected static final class PutResult<V> {
1 ✔
1583
    @Nullable V oldValue;
1584
    boolean written;
1585
  }
1586

1587
  protected enum NullCompletionListener implements CompletionListener {
1 ✔
1588
    INSTANCE;
1 ✔
1589

1590
    @Override
1591
    public void onCompletion() {}
1 ✔
1592

1593
    @Override
1594
    public void onException(Exception e) {}
1 ✔
1595
  }
1596
}
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