Skip to content

Commit 4d1c1a5

Browse files
authored
fix flakey tests (#38)
1 parent ab1f6cf commit 4d1c1a5

9 files changed

Lines changed: 841 additions & 197 deletions

File tree

core/src/main/java/net/staticstudios/data/impl/h2/H2DataAccessor.java

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -567,7 +567,15 @@ public void setRedisValue(String holderSchema, String holderTable, String identi
567567
Runnable runnable = () -> {
568568
if (value == null) {
569569
taskQueue.submitTask((connection, jedis) -> {
570-
jedis.del(key);
570+
redisListener.expectLocalDeleteEvent(key);
571+
boolean deleted = false;
572+
try {
573+
deleted = jedis.del(key) > 0;
574+
} finally {
575+
if (!deleted) {
576+
redisListener.cancelLocalDeleteEvent(key);
577+
}
578+
}
571579
});
572580
} else {
573581
taskQueue.submitTask((connection, jedis) -> {

core/src/main/java/net/staticstudios/data/impl/redis/RedisListener.java

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,13 +14,18 @@
1414
import java.util.Arrays;
1515
import java.util.Map;
1616
import java.util.Set;
17+
import java.util.concurrent.CompletableFuture;
1718
import java.util.concurrent.ConcurrentHashMap;
19+
import java.util.concurrent.TimeUnit;
20+
import java.util.concurrent.atomic.AtomicBoolean;
1821
import java.util.regex.Pattern;
1922

2023
public class RedisListener extends JedisPubSub {
2124
private static final Logger logger = LoggerFactory.getLogger(RedisListener.class);
2225
private final Set<String> listenedPartialKeys = ConcurrentHashMap.newKeySet();
2326
private final Map<Pattern, RedisEventHandler> handlers = new ConcurrentHashMap<>();
27+
private final Map<String, Integer> ignoredLocalDeleteEvents = new ConcurrentHashMap<>();
28+
private final CompletableFuture<Void> subscriptionReady = new CompletableFuture<>();
2429
private final TaskQueue taskQueue;
2530

2631
public RedisListener(DataSourceConfig ds, TaskQueue taskQueue) {
@@ -32,11 +37,18 @@ public RedisListener(DataSourceConfig ds, TaskQueue taskQueue) {
3237
if (ThreadUtils.isShuttingDown()) {
3338
return;
3439
}
40+
subscriptionReady.completeExceptionally(e);
3541
logger.error("Redis connection lost in listener thread", e);
3642
}
3743
});
3844
listenerThread.start();
3945

46+
try {
47+
subscriptionReady.get(10, TimeUnit.SECONDS);
48+
} catch (Exception e) {
49+
throw new IllegalStateException("Timed out waiting for the Redis event subscription", e);
50+
}
51+
4052
ThreadUtils.onShutdownRunSync(ShutdownStage.CLEANUP, () -> {
4153
this.punsubscribe();
4254
listenerThread.interrupt();
@@ -58,6 +70,9 @@ public void onPMessage(String pattern, String channel, String key) {
5870
if (!key.startsWith("static-data:")) {
5971
return;
6072
}
73+
if (event == RedisEvent.DEL && consumeLocalDeleteEvent(key)) {
74+
return;
75+
}
6176

6277
for (Map.Entry<Pattern, RedisEventHandler> entry : handlers.entrySet()) {
6378
if (entry.getKey().matcher(key).matches()) {
@@ -75,4 +90,26 @@ public void onPMessage(String pattern, String channel, String key) {
7590
}
7691
}
7792
}
93+
94+
@Override
95+
public void onPSubscribe(String pattern, int subscribedChannels) {
96+
subscriptionReady.complete(null);
97+
}
98+
99+
public void expectLocalDeleteEvent(String key) {
100+
ignoredLocalDeleteEvents.merge(key, 1, Integer::sum);
101+
}
102+
103+
public void cancelLocalDeleteEvent(String key) {
104+
consumeLocalDeleteEvent(key);
105+
}
106+
107+
private boolean consumeLocalDeleteEvent(String key) {
108+
AtomicBoolean consumed = new AtomicBoolean(false);
109+
ignoredLocalDeleteEvents.computeIfPresent(key, (ignoredKey, count) -> {
110+
consumed.set(true);
111+
return count == 1 ? null : count - 1;
112+
});
113+
return consumed.get();
114+
}
78115
}

core/src/test/java/net/staticstudios/data/CachedValueTest.java

Lines changed: 21 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@
1111
import org.junit.jupiter.api.Test;
1212
import redis.clients.jedis.Jedis;
1313

14+
import java.util.Objects;
1415
import java.util.UUID;
1516

1617
import static org.junit.jupiter.api.Assertions.*;
@@ -52,7 +53,7 @@ public void testFallback() {
5253
assertEquals(false, user.onCooldown.get());
5354
assertEquals(0, user.cooldownUpdates.get());
5455

55-
waitForDataPropagation();
56+
flushDataManagers();
5657

5758
Jedis jedis = getJedis();
5859

@@ -114,9 +115,13 @@ public void testUpdateHandler() {
114115

115116
Jedis jedis = getJedis();
116117
String onCooldownKey = RedisUtils.buildRedisKey("public", "users", "on_cooldown", user.getIdColumns());
118+
dataManager.flushTaskQueue();
117119
jedis.del(onCooldownKey);
118120

119-
waitForDataPropagation();
121+
awaitCondition(
122+
() -> Objects.equals(false, user.onCooldown.get()) && Objects.equals(6, user.cooldownUpdates.get()),
123+
"the external Redis deletion to reach the cached value and its update handler"
124+
);
120125

121126
assertEquals(false, user.onCooldown.get());
122127
assertEquals(6, user.cooldownUpdates.get());
@@ -139,19 +144,19 @@ public void testUpdateRedis() {
139144

140145
user.onCooldown.set(true);
141146
user.cooldownUpdates.set(1);
142-
waitForDataPropagation();
147+
dataManager.flushTaskQueue();
143148
assertEquals("true", gson.fromJson(jedis.get(onCooldownKey), RedisEncodedValue.class).value());
144149
assertEquals("1", gson.fromJson(jedis.get(cooldownUpdatesKey), RedisEncodedValue.class).value());
145150

146151
user.onCooldown.set(null);
147152
user.cooldownUpdates.set(null);
148-
waitForDataPropagation();
153+
dataManager.flushTaskQueue();
149154
assertNull(jedis.get(onCooldownKey));
150155
assertNull(jedis.get(cooldownUpdatesKey));
151156

152157
user.onCooldown.set(false); //fallback
153158
user.cooldownUpdates.set(0); //fallback
154-
waitForDataPropagation();
159+
dataManager.flushTaskQueue();
155160
assertNull(jedis.get(onCooldownKey));
156161
assertNull(jedis.get(cooldownUpdatesKey));
157162
}
@@ -175,7 +180,7 @@ public void testLoadCachedValues() {
175180

176181
jedis.set(cooldownUpdatesKey, gson.toJson(new RedisEncodedValue(null, "5")));
177182

178-
waitForDataPropagation();
183+
awaitCondition(() -> Objects.equals(5, user1.cooldownUpdates.get()), "the external Redis value to reach H2");
179184

180185
assertEquals(5, user1.cooldownUpdates.get());
181186

@@ -204,15 +209,15 @@ public void testRefreshCachedValues() throws InterruptedException {
204209
assertEquals(2, user.counter.refresh());
205210
assertEquals(2, user.counter.get());
206211

207-
Thread.sleep(10_000); //wait for the cached value to expire
208-
209212
String counterKey = RedisUtils.buildRedisKey("public", "users", "counter", user.getIdColumns());
210213

211214
Jedis jedis = getJedis();
215+
awaitCondition(() -> !jedis.exists(counterKey), "the cached counter to expire");
212216
assertFalse(jedis.exists(counterKey));
213217

214-
assertEquals(0, user.counter.get()); //trigger a refresh
215-
waitForDataPropagation();
218+
awaitCondition(() -> Objects.equals(0, user.counter.get()), "the expiration event to clear and refresh the H2 cached value");
219+
assertEquals(0, user.counter.get());
220+
dataManager.flushTaskQueue();
216221
assertEquals("0", gson.fromJson(jedis.get(counterKey), RedisEncodedValue.class).value());
217222
}
218223

@@ -234,17 +239,20 @@ public void testUpdateInterval() throws Exception {
234239
}
235240

236241
assertEquals(4, user.throttledCounter.get());
237-
waitForDataPropagation();
242+
dataManager.flushTaskQueue();
238243

239244
Jedis jedis = getJedis();
240245
String throttledCounterKey = RedisUtils.buildRedisKey("public", "users", "throttled_counter", user.getIdColumns());
241246

242247
assertNull(jedis.get(throttledCounterKey));
243248

244-
Thread.sleep(6000);
249+
awaitCondition(() -> {
250+
RedisEncodedValue value = gson.fromJson(jedis.get(throttledCounterKey), RedisEncodedValue.class);
251+
return value != null && Objects.equals("4", value.value());
252+
}, "the throttled cached value to be written");
245253

246254
RedisEncodedValue encoded = gson.fromJson(jedis.get(throttledCounterKey), RedisEncodedValue.class);
247255
assertNotNull(encoded);
248256
assertEquals("4", encoded.value());
249257
}
250-
}
258+
}

core/src/test/java/net/staticstudios/data/PersistentManyToManyCollectionTest.java

Lines changed: 15 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -319,7 +319,7 @@ public void testAddHandlerUpdate() throws SQLException {
319319
.insert(InsertMode.SYNC);
320320
List<MockUser> friends = createFriends(5);
321321
user.friends.addAll(friends);
322-
waitForDataPropagation();
322+
flushDataManagers();
323323
assertEquals(5, user.friendAdditions.get());
324324

325325
List<MockUser> otherFriends = createFriends(5);
@@ -334,8 +334,12 @@ public void testAddHandlerUpdate() throws SQLException {
334334
preparedStatement.setObject(3, friend.id.get());
335335
preparedStatement.executeUpdate();
336336
}
337-
waitForDataPropagation();
338-
assertEquals(5 + (++i), user.friendAdditions.get());
337+
int expectedAdditions = 5 + (++i);
338+
awaitCondition(
339+
() -> user.friendAdditions.get() == expectedAdditions,
340+
"the many-to-many update addition handler"
341+
);
342+
assertEquals(expectedAdditions, user.friendAdditions.get());
339343
}
340344
}
341345

@@ -347,7 +351,7 @@ public void testRemoveHandlerUpdate() throws SQLException {
347351
.insert(InsertMode.SYNC);
348352
List<MockUser> friends = createFriends(5);
349353
user.friends.addAll(friends);
350-
waitForDataPropagation();
354+
flushDataManagers();
351355
assertEquals(5, user.friendAdditions.get());
352356

353357
Connection pgConnection = getConnection();
@@ -359,8 +363,12 @@ public void testRemoveHandlerUpdate() throws SQLException {
359363
preparedStatement.setObject(2, friend.id.get());
360364
preparedStatement.executeUpdate();
361365
}
362-
waitForDataPropagation();
363-
assertEquals(++i, user.friendRemovals.get());
366+
int expectedRemovals = ++i;
367+
awaitCondition(
368+
() -> user.friendRemovals.get() == expectedRemovals,
369+
"the many-to-many delete handler"
370+
);
371+
assertEquals(expectedRemovals, user.friendRemovals.get());
364372
}
365373
}
366374

@@ -402,4 +410,4 @@ public void testRemoveHandlerDelete() {
402410
assertEquals(++i, user.friendRemovals.get());
403411
}
404412
}
405-
}
413+
}

core/src/test/java/net/staticstudios/data/PersistentOneToManyValueCollectionTest.java

Lines changed: 13 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -337,9 +337,12 @@ public void testAddHandlerUpdate() {
337337
} catch (Exception e) {
338338
throw new RuntimeException(e);
339339
}
340-
waitForDataPropagation();
341-
342-
assertEquals(++i, user.favoriteNumberAdditions.get());
340+
int expectedAdditions = ++i;
341+
awaitCondition(
342+
() -> user.favoriteNumberAdditions.get() == expectedAdditions,
343+
"the one-to-many value update addition handler"
344+
);
345+
assertEquals(expectedAdditions, user.favoriteNumberAdditions.get());
343346
}
344347
}
345348

@@ -381,9 +384,12 @@ public void testRemoveHandlerUpdate() {
381384
} catch (Exception e) {
382385
throw new RuntimeException(e);
383386
}
384-
waitForDataPropagation();
385-
386-
assertEquals(++i, user.favoriteNumberRemovals.get());
387+
int expectedRemovals = ++i;
388+
awaitCondition(
389+
() -> user.favoriteNumberRemovals.get() == expectedRemovals,
390+
"the one-to-many value update removal handler"
391+
);
392+
assertEquals(expectedRemovals, user.favoriteNumberRemovals.get());
387393
}
388394
}
389395

@@ -425,4 +431,4 @@ public void testRemoveHandlerDelete() {
425431
assertEquals(++i, user.favoriteNumberRemovals.get());
426432
}
427433
}
428-
}
434+
}

0 commit comments

Comments
 (0)