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

openmrs / openmrs-core / 37033370540

02 Oct 2026 04:21PM UTC coverage: 66.43% (+0.1%) from 66.293%
37033370540

push

github

ibacher
TRUNK-6810: Extend cluster-safety to RolePrivilegeCache (#6627)

281 of 297 new or added lines in 9 files covered. (94.61%)

2 existing lines in 1 file now uncovered.

25927 of 39029 relevant lines covered (66.43%)

0.66 hits per line

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

97.47
/api/src/main/java/org/openmrs/api/cache/CacheInvalidation.java
1
/**
2
 * This Source Code Form is subject to the terms of the Mozilla Public License,
3
 * v. 2.0. If a copy of the MPL was not distributed with this file, You can
4
 * obtain one at http://mozilla.org/MPL/2.0/. OpenMRS is also distributed under
5
 * the terms of the Healthcare Disclaimer located at http://openmrs.org/license.
6
 *
7
 * Copyright (C) OpenMRS Inc. OpenMRS is a registered trademark and the OpenMRS
8
 * graphic logo is a trademark of OpenMRS Inc.
9
 */
10
package org.openmrs.api.cache;
11

12
import java.util.Collections;
13
import java.util.HashSet;
14
import java.util.Set;
15
import java.util.UUID;
16
import java.util.concurrent.CompletableFuture;
17
import java.util.function.Supplier;
18

19
import org.infinispan.Cache;
20
import org.infinispan.context.Flag;
21
import org.infinispan.manager.EmbeddedCacheManager;
22
import org.infinispan.notifications.Listener;
23
import org.infinispan.notifications.cachemanagerlistener.annotation.Merged;
24
import org.infinispan.notifications.cachemanagerlistener.event.MergeEvent;
25
import org.slf4j.Logger;
26
import org.slf4j.LoggerFactory;
27
import org.springframework.transaction.support.TransactionSynchronization;
28
import org.springframework.transaction.support.TransactionSynchronizationManager;
29

30
/**
31
 * Invalidation for an API cache whose values are loaded from committed data, possibly on another
32
 * thread, while the cache may be shared by a cluster. Each cache owns one instance.
33
 * <p>
34
 * Every eviction replaces a generation token kept in the cache itself. A load reads the token with
35
 * {@link #currentGeneration(Cache)} before reading the database and writes its value with
36
 * {@link #putIfCurrent(Cache, Object, Object, Object)}, which only keeps it if nothing was evicted
37
 * in between. In a cluster, replacing the token invalidates it on every node, so an eviction
38
 * anywhere is seen everywhere. A token that expires or is evicted for space also reads as an
39
 * eviction, which only costs an uncached load. If an eviction cannot reach every node, this node
40
 * clears its own cache instead, and a node clears its cache when a network partition that cut it
41
 * off heals.
42
 * <p>
43
 * Within a transaction, {@link #invalidate(String)} records the key rather than evicting it, and
44
 * the keys a transaction records are evicted together when it completes, whether it commits or
45
 * rolls back. Until then {@link #isWrittenInCurrentTransaction(String)} reports them, so the cache
46
 * can let the transaction read its own writes without the cache while other transactions keep
47
 * reading the committed values. Outside a transaction keys are evicted immediately.
48
 *
49
 * @since 2.8.10
50
 */
51
public final class CacheInvalidation {
52

53
        private static final Logger log = LoggerFactory.getLogger(CacheInvalidation.class);
1 ✔
54

55
        /**
56
         * Key of the generation token. The NUL character keeps it from colliding with a cache's own keys.
57
         */
58
        static final String GENERATION = "\0generation";
59

60
        /** Recorded by {@link #invalidate(String)} to mean every key. */
61
        static final String ALL = "\0all";
62

63
        private final String cacheDescription;
64

65
        private final Supplier<Cache<Object, Object>> cache;
66

67
        private final EmbeddedCacheManager cacheManager;
68

69
        private final PartitionMergeListener mergeListener = new PartitionMergeListener(this);
1 ✔
70

71
        /**
72
         * Creates the invalidation for a cache and starts listening for partition merges; call
73
         * {@link #close()} when the cache is destroyed.
74
         *
75
         * @param cacheDescription how log messages name the cache, for example "the role privilege cache"
76
         * @param cache supplies the cache, or null if it is unavailable
77
         * @param cacheManager the cache manager the cache belongs to
78
         */
79
        CacheInvalidation(String cacheDescription, Supplier<Cache<Object, Object>> cache, EmbeddedCacheManager cacheManager) {
1 ✔
80
                this.cacheDescription = cacheDescription;
1 ✔
81
                this.cache = cache;
1 ✔
82
                this.cacheManager = cacheManager;
1 ✔
83
                cacheManager.addListener(mergeListener);
1 ✔
84
        }
1 ✔
85

86
        /** Stops listening for partition merges. */
87
        void close() {
88
                cacheManager.removeListener(mergeListener);
1 ✔
89
        }
1 ✔
90

91
        /**
92
         * Evicts <code>key</code>, or every key if it is {@link #ALL}, when the current transaction
93
         * completes, or immediately outside one.
94
         *
95
         * @param key the key to evict
96
         */
97
        void invalidate(String key) {
98
                Cache<Object, Object> current = cache.get();
1 ✔
99
                if (current == null) {
1 ✔
NEW
100
                        return;
×
101
                }
102

103
                if (TransactionSynchronizationManager.isSynchronizationActive()) {
1 ✔
104
                        getWrittenInCurrentTransaction(current).add(key);
1 ✔
105
                } else {
106
                        evictOrClearLocally(current, Collections.singleton(key));
1 ✔
107
                }
108
        }
1 ✔
109

110
        /**
111
         * @param key the key to check, or null to check only whether every key has been invalidated
112
         * @return true if the current transaction has invalidated <code>key</code> or every key
113
         */
114
        boolean isWrittenInCurrentTransaction(String key) {
115
                Set<?> written = (Set<?>) TransactionSynchronizationManager.getResource(this);
1 ✔
116
                return written != null && (written.contains(ALL) || (key != null && written.contains(key)));
1 ✔
117
        }
118

119
        /**
120
         * Evicts every key immediately, even within a transaction, which still reads without the cache
121
         * until it completes. Only for tests that load data behind the API.
122
         */
123
        void clearNow() {
124
                Cache<Object, Object> current = cache.get();
1 ✔
125
                if (current == null) {
1 ✔
NEW
126
                        return;
×
127
                }
128

129
                evictOrClearLocally(current, Collections.singleton(ALL));
1 ✔
130
                if (TransactionSynchronizationManager.isSynchronizationActive()) {
1 ✔
131
                        getWrittenInCurrentTransaction(current).add(ALL);
1 ✔
132
                }
133
        }
1 ✔
134

135
        /**
136
         * Forgets which keys the current transaction has invalidated, so that it reads through the cache
137
         * again and does not evict them when it completes. Only for tests that load data behind the API but
138
         * still need to observe cache hits.
139
         */
140
        void forgetWritesInCurrentTransaction() {
141
                Set<?> written = (Set<?>) TransactionSynchronizationManager.getResource(this);
1 ✔
142
                if (written != null) {
1 ✔
143
                        written.clear();
1 ✔
144
                }
145
        }
1 ✔
146

147
        /**
148
         * Returns this node's generation token, first creating one if there is none. The token is created
149
         * on this node only, since writing it cluster-wide would invalidate the other nodes' tokens.
150
         *
151
         * @param cache the cache
152
         * @return the token to pass to {@link #putIfCurrent(Cache, Object, Object, Object)}
153
         */
154
        static Object currentGeneration(Cache<Object, Object> cache) {
155
                Object generation = cache.get(GENERATION);
1 ✔
156
                if (generation != null) {
1 ✔
157
                        return generation;
1 ✔
158
                }
159

160
                String created = UUID.randomUUID().toString();
1 ✔
161
                Object existing = cache.getAdvancedCache().withFlags(Flag.CACHE_MODE_LOCAL).putIfAbsent(GENERATION, created);
1 ✔
162
                return existing != null ? existing : created;
1 ✔
163
        }
164

165
        /**
166
         * Caches <code>value</code> unless an eviction has happened since <code>generation</code> was read.
167
         * An eviction replaces the generation token before removing entries, so if one runs while the value
168
         * is being put, the second check either sees it and removes the value, or the eviction removes it.
169
         * The value is written with {@code putForExternalRead}, so it does not invalidate other nodes'
170
         * entries, and the removal is local, since the value was only ever written on this node.
171
         *
172
         * @param cache the cache
173
         * @param key the key to cache the value under
174
         * @param value the value
175
         * @param generation the token {@link #currentGeneration(Cache)} returned before the value was read
176
         */
177
        static void putIfCurrent(Cache<Object, Object> cache, Object key, Object value, Object generation) {
178
                if (!isCurrent(cache, generation)) {
1 ✔
179
                        return;
1 ✔
180
                }
181

182
                cache.putForExternalRead(key, value);
1 ✔
183
                if (!isCurrent(cache, generation)) {
1 ✔
184
                        cache.getAdvancedCache().withFlags(Flag.CACHE_MODE_LOCAL).remove(key);
1 ✔
185
                }
186
        }
1 ✔
187

188
        /**
189
         * @param cache the cache
190
         * @param generation a token {@link #currentGeneration(Cache)} returned, or null
191
         * @return true if nothing has been evicted on this node since <code>generation</code> was read
192
         */
193
        static boolean isCurrent(Cache<Object, Object> cache, Object generation) {
194
                return generation != null && generation.equals(cache.get(GENERATION));
1 ✔
195
        }
196

197
        /**
198
         * Evicts the keys on every node or, if that fails, for example because the cluster is partitioned,
199
         * clears this node's cache instead, so that at least this node does not serve values from before
200
         * the write. Nodes the eviction did not reach may serve them until the lifespan expires, or, if
201
         * they were cut off by a partition, until it heals.
202
         */
203
        private void evictOrClearLocally(Cache<Object, Object> cache, Set<String> keys) {
204
                try {
205
                        evictNow(cache, keys);
1 ✔
206
                } catch (RuntimeException e) {
1 ✔
207
                        log.error("Could not evict {} from {} on every node, so clearing it on this node only",
1 ✔
208
                            keys.contains(ALL) ? "every entry" : keys, cacheDescription, e);
1 ✔
209
                        clearLocally(cache);
1 ✔
210
                }
1 ✔
211
        }
1 ✔
212

213
        /**
214
         * Replaces the generation token and then removes the keys, or every entry if they include
215
         * {@link #ALL}. The token must be replaced first: a clear does not lock every key, so a load could
216
         * otherwise write behind it and still find its token. Replacing the token with a plain put
217
         * invalidates it on every other node before the removals are sent. The removals are sent together,
218
         * so evicting several keys costs about one round trip to the other nodes.
219
         */
220
        private static void evictNow(Cache<Object, Object> cache, Set<String> keys) {
221
                cache.put(GENERATION, UUID.randomUUID().toString());
1 ✔
222
                if (keys.contains(ALL)) {
1 ✔
223
                        cache.clear();
1 ✔
224
                } else {
225
                        CompletableFuture.allOf(keys.stream().map(cache::removeAsync).toArray(CompletableFuture[]::new)).join();
1 ✔
226
                }
227
        }
1 ✔
228

229
        /**
230
         * Clears this node's entries, including its generation token, so that loads in progress on this
231
         * node are discarded too.
232
         */
233
        private static void clearLocally(Cache<Object, Object> cache) {
234
                cache.getAdvancedCache().withFlags(Flag.CACHE_MODE_LOCAL).clear();
1 ✔
235
        }
1 ✔
236

237
        /**
238
         * Returns the keys the current transaction has invalidated, first arranging for them to be evicted
239
         * when it completes. The set is unbound while the transaction is suspended, since Spring only
240
         * suspends its own resources, so a transaction started in the meantime, for example with
241
         * {@code REQUIRES_NEW}, records and evicts its own writes.
242
         */
243
        @SuppressWarnings("unchecked")
244
        private Set<String> getWrittenInCurrentTransaction(Cache<Object, Object> cache) {
245
                Set<String> bound = (Set<String>) TransactionSynchronizationManager.getResource(this);
1 ✔
246
                if (bound != null) {
1 ✔
247
                        return bound;
1 ✔
248
                }
249

250
                Set<String> written = new HashSet<>();
1 ✔
251
                TransactionSynchronizationManager.bindResource(this, written);
1 ✔
252
                TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() {
1 ✔
253

254
                        @Override
255
                        public void suspend() {
256
                                TransactionSynchronizationManager.unbindResource(CacheInvalidation.this);
1 ✔
257
                        }
1 ✔
258

259
                        @Override
260
                        public void resume() {
261
                                TransactionSynchronizationManager.bindResource(CacheInvalidation.this, written);
1 ✔
262
                        }
1 ✔
263

264
                        @Override
265
                        public void afterCompletion(int status) {
266
                                try {
267
                                        if (!written.isEmpty()) {
1 ✔
268
                                                evictOrClearLocally(cache, written);
1 ✔
269
                                        }
270
                                } finally {
271
                                        TransactionSynchronizationManager.unbindResourceIfPossible(CacheInvalidation.this);
1 ✔
272
                                }
273
                        }
1 ✔
274
                });
275
                return written;
1 ✔
276
        }
277

278
        /**
279
         * Clears this node's cache when a network partition heals, since while it was cut off it missed the
280
         * evictions of writes made on the other side. Public only because Infinispan requires listeners to
281
         * be.
282
         */
283
        @Listener
284
        public static final class PartitionMergeListener {
285

286
                private final CacheInvalidation owner;
287

288
                PartitionMergeListener(CacheInvalidation owner) {
1 ✔
289
                        this.owner = owner;
1 ✔
290
                }
1 ✔
291

292
                @Merged
293
                public void merged(MergeEvent event) {
294
                        Cache<Object, Object> cache = owner.cache.get();
1 ✔
295
                        if (cache != null) {
1 ✔
296
                                clearLocally(cache);
1 ✔
297
                        }
298
                }
1 ✔
299
        }
300
}
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