1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17 package org.apache.commons.lang3.concurrent;
18
19 import static org.apache.commons.lang3.LangAssertions.assertIllegalArgumentException;
20 import static org.junit.jupiter.api.Assertions.assertEquals;
21 import static org.junit.jupiter.api.Assertions.assertFalse;
22 import static org.junit.jupiter.api.Assertions.assertNotNull;
23 import static org.junit.jupiter.api.Assertions.assertThrows;
24 import static org.junit.jupiter.api.Assertions.assertTimeoutPreemptively;
25 import static org.junit.jupiter.api.Assertions.assertTrue;
26
27 import java.time.Duration;
28 import java.util.concurrent.CountDownLatch;
29 import java.util.concurrent.ScheduledExecutorService;
30 import java.util.concurrent.ScheduledFuture;
31 import java.util.concurrent.ScheduledThreadPoolExecutor;
32 import java.util.concurrent.TimeUnit;
33
34 import org.apache.commons.lang3.AbstractLangTest;
35 import org.apache.commons.lang3.ThreadUtils;
36 import org.easymock.EasyMock;
37 import org.junit.jupiter.api.Test;
38
39
40
41
42 class TimedSemaphoreTest extends AbstractLangTest {
43
44
45
46
47
48
49 private static final class SemaphoreThread extends Thread {
50
51
52 private final TimedSemaphore semaphore;
53
54
55 private final CountDownLatch latch;
56
57
58 private final int count;
59
60
61 private final int latchCount;
62
63 SemaphoreThread(final TimedSemaphore b, final CountDownLatch l, final int c, final int lc) {
64 semaphore = b;
65 latch = l;
66 count = c;
67 latchCount = lc;
68 }
69
70
71
72
73
74
75 @Override
76 public void run() {
77 try {
78 for (int i = 0; i < count; i++) {
79 semaphore.acquire();
80
81 if (i < latchCount) {
82 latch.countDown();
83 }
84 }
85 } catch (final InterruptedException iex) {
86 Thread.currentThread().interrupt();
87 }
88 }
89 }
90
91
92
93
94
95 private static final class TimedSemaphoreTestImpl extends TimedSemaphore {
96
97
98 ScheduledFuture<?> schedFuture;
99
100
101 volatile CountDownLatch latch;
102
103
104 private int periodEnds;
105
106 TimedSemaphoreTestImpl(final long timePeriod, final TimeUnit timeUnit, final int limit) {
107 super(timePeriod, timeUnit, limit);
108 }
109
110 TimedSemaphoreTestImpl(final ScheduledExecutorService service, final long timePeriod, final TimeUnit timeUnit, final int limit) {
111 super(service, timePeriod, timeUnit, limit);
112 }
113
114
115
116
117
118
119 @Override
120 public synchronized void acquire() throws InterruptedException {
121 super.acquire();
122 if (latch != null) {
123 latch.countDown();
124 }
125 }
126
127
128
129
130 @Override
131 protected synchronized void endOfPeriod() {
132 super.endOfPeriod();
133 periodEnds++;
134 }
135
136
137
138
139
140
141 int getPeriodEnds() {
142 synchronized (this) {
143 return periodEnds;
144 }
145 }
146
147
148
149
150 @Override
151 protected ScheduledFuture<?> startTimer() {
152 return schedFuture != null ? schedFuture : super.startTimer();
153 }
154 }
155
156
157
158
159
160 private static final class TryAcquireThread extends Thread {
161
162
163 private final TimedSemaphore semaphore;
164
165
166 private final CountDownLatch latch;
167
168
169 private boolean acquired;
170
171 TryAcquireThread(final TimedSemaphore s, final CountDownLatch l) {
172 semaphore = s;
173 latch = l;
174 }
175
176 @Override
177 public void run() {
178 try {
179 if (latch.await(10, TimeUnit.SECONDS)) {
180 acquired = semaphore.tryAcquire();
181 }
182 } catch (final InterruptedException iex) {
183
184 }
185 }
186 }
187
188
189 private static final long PERIOD_MILLIS = 500;
190
191 private static final Duration DURATION = Duration.ofMillis(PERIOD_MILLIS);
192
193
194 private static final TimeUnit UNIT = TimeUnit.MILLISECONDS;
195
196
197 private static final int LIMIT = 10;
198
199
200
201
202
203
204
205 private void prepareStartTimer(final ScheduledExecutorService service,
206 final ScheduledFuture<?> future) {
207 service.scheduleAtFixedRate((Runnable) EasyMock.anyObject(), EasyMock.eq(PERIOD_MILLIS), EasyMock.eq(PERIOD_MILLIS), EasyMock.eq(UNIT));
208 EasyMock.expectLastCall().andReturn(future);
209 }
210
211
212
213
214
215
216 @Test
217 void testAcquireLimit() throws InterruptedException {
218 final ScheduledExecutorService service = EasyMock.createMock(ScheduledExecutorService.class);
219 final ScheduledFuture<?> future = EasyMock.createMock(ScheduledFuture.class);
220 prepareStartTimer(service, future);
221 EasyMock.replay(service, future);
222 final int count = 10;
223 final CountDownLatch latch = new CountDownLatch(count - 1);
224 final TimedSemaphore semaphore = new TimedSemaphore(service, PERIOD_MILLIS, UNIT, 1);
225 final SemaphoreThread t = new SemaphoreThread(semaphore, latch, count, count - 1);
226 semaphore.setLimit(count - 1);
227
228 t.start();
229 latch.await();
230
231 assertEquals(count - 1, semaphore.getAcquireCount(), "Wrong semaphore count");
232
233 semaphore.endOfPeriod();
234 t.join();
235 assertEquals(1, semaphore.getAcquireCount(), "Wrong semaphore count (2)");
236 assertEquals(count - 1, semaphore.getLastAcquiresPerPeriod(), "Wrong acquire() count");
237 EasyMock.verify(service, future);
238 }
239
240
241
242
243
244
245
246
247
248 @Test
249 void testAcquireMultiplePeriods() throws InterruptedException {
250 final int count = 1000;
251 final TimedSemaphoreTestImpl semaphore = new TimedSemaphoreTestImpl(PERIOD_MILLIS / 10, TimeUnit.MILLISECONDS, 1);
252 semaphore.setLimit(count / 4);
253 final CountDownLatch latch = new CountDownLatch(count);
254 final SemaphoreThread t = new SemaphoreThread(semaphore, latch, count, count);
255 t.start();
256 latch.await();
257 semaphore.shutdown();
258 assertTrue(semaphore.getPeriodEnds() > 0, "End of period not reached");
259 }
260
261
262
263
264
265
266
267
268
269 @Test
270 void testAcquireMultipleThreads() throws InterruptedException {
271 final ScheduledExecutorService service = EasyMock.createMock(ScheduledExecutorService.class);
272 final ScheduledFuture<?> future = EasyMock.createMock(ScheduledFuture.class);
273 prepareStartTimer(service, future);
274 EasyMock.replay(service, future);
275 final TimedSemaphoreTestImpl semaphore = new TimedSemaphoreTestImpl(service, PERIOD_MILLIS, UNIT, 1);
276 semaphore.latch = new CountDownLatch(1);
277 final int count = 10;
278 final SemaphoreThread[] threads = new SemaphoreThread[count];
279 for (int i = 0; i < count; i++) {
280 threads[i] = new SemaphoreThread(semaphore, null, 1, 0);
281 threads[i].start();
282 }
283 for (int i = 0; i < count; i++) {
284 semaphore.latch.await();
285 assertEquals(1, semaphore.getAcquireCount(), "Wrong count");
286 semaphore.latch = new CountDownLatch(1);
287 semaphore.endOfPeriod();
288 assertEquals(1, semaphore.getLastAcquiresPerPeriod(), "Wrong acquire count");
289 }
290 for (int i = 0; i < count; i++) {
291 threads[i].join();
292 }
293 EasyMock.verify(service, future);
294 }
295
296
297
298
299
300
301
302
303 @Test
304 void testAcquireNoLimit() throws InterruptedException {
305 final ScheduledExecutorService service = EasyMock.createMock(ScheduledExecutorService.class);
306 final ScheduledFuture<?> future = EasyMock.createMock(ScheduledFuture.class);
307 prepareStartTimer(service, future);
308 EasyMock.replay(service, future);
309 final TimedSemaphoreTestImpl semaphore = new TimedSemaphoreTestImpl(service, PERIOD_MILLIS, UNIT, TimedSemaphore.NO_LIMIT);
310 final int count = 1000;
311 final CountDownLatch latch = new CountDownLatch(count);
312 final SemaphoreThread t = new SemaphoreThread(semaphore, latch, count, count);
313 t.start();
314 latch.await();
315 EasyMock.verify(service, future);
316 }
317
318
319
320
321
322
323 @Test
324 void testGetAvailablePermits() throws InterruptedException {
325 final ScheduledExecutorService service = EasyMock.createMock(ScheduledExecutorService.class);
326 final ScheduledFuture<?> future = EasyMock.createMock(ScheduledFuture.class);
327 prepareStartTimer(service, future);
328 EasyMock.replay(service, future);
329 final TimedSemaphore semaphore = new TimedSemaphore(service, PERIOD_MILLIS, UNIT, LIMIT);
330 for (int i = 0; i < LIMIT; i++) {
331 assertEquals(LIMIT - i, semaphore.getAvailablePermits(), "Wrong available count at " + i);
332 semaphore.acquire();
333 }
334 semaphore.endOfPeriod();
335 assertEquals(LIMIT, semaphore.getAvailablePermits(), "Wrong available count in new period");
336 EasyMock.verify(service, future);
337 }
338
339
340
341
342
343
344 @Test
345 void testGetAverageCallsPerPeriod() throws InterruptedException {
346 final ScheduledExecutorService service = EasyMock.createMock(ScheduledExecutorService.class);
347 final ScheduledFuture<?> future = EasyMock.createMock(ScheduledFuture.class);
348 prepareStartTimer(service, future);
349 EasyMock.replay(service, future);
350 final TimedSemaphore semaphore = new TimedSemaphore(service, PERIOD_MILLIS, UNIT, LIMIT);
351 semaphore.acquire();
352 semaphore.endOfPeriod();
353 assertEquals(1.0, semaphore.getAverageCallsPerPeriod(), .005, "Wrong average (1)");
354 semaphore.acquire();
355 semaphore.acquire();
356 semaphore.endOfPeriod();
357 assertEquals(1.5, semaphore.getAverageCallsPerPeriod(), .005, "Wrong average (2)");
358 EasyMock.verify(service, future);
359 }
360
361
362
363
364 @Test
365 void testInit() {
366 final ScheduledExecutorService service = EasyMock.createMock(ScheduledExecutorService.class);
367 EasyMock.replay(service);
368 final TimedSemaphore semaphore = new TimedSemaphore(service, PERIOD_MILLIS, UNIT, LIMIT);
369 EasyMock.verify(service);
370 assertEquals(service, semaphore.getExecutorService(), "Wrong service");
371 assertEquals(PERIOD_MILLIS, semaphore.getPeriod(), "Wrong period");
372 assertEquals(UNIT, semaphore.getUnit(), "Wrong unit");
373 assertEquals(0, semaphore.getLastAcquiresPerPeriod(), "Statistic available");
374 assertEquals(0.0, semaphore.getAverageCallsPerPeriod(), .05, "Average available");
375 assertFalse(semaphore.isShutdown(), "Already shutdown");
376 assertEquals(LIMIT, semaphore.getLimit(), "Wrong limit");
377 }
378
379
380
381
382
383 @Test
384 void testInitDefaultService() {
385 final TimedSemaphore semaphore = new TimedSemaphore(PERIOD_MILLIS, UNIT, LIMIT);
386 final ScheduledThreadPoolExecutor exec = (ScheduledThreadPoolExecutor) semaphore.getExecutorService();
387 assertFalse(exec.getContinueExistingPeriodicTasksAfterShutdownPolicy(), "Wrong periodic task policy");
388 assertFalse(exec.getExecuteExistingDelayedTasksAfterShutdownPolicy(), "Wrong delayed task policy");
389 assertFalse(exec.isShutdown(), "Already shutdown");
390 semaphore.shutdown();
391 }
392
393
394
395
396
397 @Test
398 void testInitInvalidPeriod() {
399 assertIllegalArgumentException(() -> new TimedSemaphore(0L, UNIT, LIMIT));
400 }
401
402
403
404
405 @Test
406 void testPassAfterShutdown() {
407 final TimedSemaphore semaphore = new TimedSemaphore(PERIOD_MILLIS, UNIT, LIMIT);
408 semaphore.shutdown();
409 assertThrows(IllegalStateException.class, semaphore::acquire);
410 }
411
412
413
414
415
416
417 @Test
418 void testShutdownMultipleTimes() throws InterruptedException {
419 final ScheduledExecutorService service = EasyMock.createMock(ScheduledExecutorService.class);
420 final ScheduledFuture<?> future = EasyMock.createMock(ScheduledFuture.class);
421 prepareStartTimer(service, future);
422 EasyMock.expect(Boolean.valueOf(future.cancel(false))).andReturn(Boolean.TRUE);
423 EasyMock.replay(service, future);
424 final TimedSemaphoreTestImpl semaphore = new TimedSemaphoreTestImpl(service, PERIOD_MILLIS, UNIT, LIMIT);
425 semaphore.acquire();
426 for (int i = 0; i < 10; i++) {
427 semaphore.shutdown();
428 }
429 EasyMock.verify(service, future);
430 }
431
432
433
434
435
436 @Test
437 void testShutdownOwnExecutor() {
438 final TimedSemaphore semaphore = new TimedSemaphore(PERIOD_MILLIS, UNIT, LIMIT);
439 semaphore.shutdown();
440 assertTrue(semaphore.isShutdown(), "Not shutdown");
441 assertTrue(semaphore.getExecutorService().isShutdown(), "Executor not shutdown");
442 }
443
444
445
446
447
448 @Test
449 void testShutdownSharedExecutorNoTask() {
450 final ScheduledExecutorService service = EasyMock.createMock(ScheduledExecutorService.class);
451 EasyMock.replay(service);
452 final TimedSemaphore semaphore = new TimedSemaphore(service, PERIOD_MILLIS, UNIT, LIMIT);
453 semaphore.shutdown();
454 assertTrue(semaphore.isShutdown(), "Not shutdown");
455 EasyMock.verify(service);
456 }
457
458
459
460
461
462
463
464 @Test
465 void testShutdownSharedExecutorTask() throws InterruptedException {
466 final ScheduledExecutorService service = EasyMock.createMock(ScheduledExecutorService.class);
467 final ScheduledFuture<?> future = EasyMock.createMock(ScheduledFuture.class);
468 prepareStartTimer(service, future);
469 EasyMock.expect(Boolean.valueOf(future.cancel(false))).andReturn(Boolean.TRUE);
470 EasyMock.replay(service, future);
471 final TimedSemaphoreTestImpl semaphore = new TimedSemaphoreTestImpl(service, PERIOD_MILLIS, UNIT, LIMIT);
472 semaphore.acquire();
473 semaphore.shutdown();
474 assertTrue(semaphore.isShutdown(), "Not shutdown");
475 EasyMock.verify(service, future);
476 }
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504 @Test
505 public void testShutdownWakesBlockedAcquireThreads() {
506 assertTimeoutPreemptively(Duration.ofSeconds(10), () -> {
507
508
509
510 final TimedSemaphore sem = TimedSemaphore.builder().setPeriod(60).setTimeUnit(TimeUnit.SECONDS).setLimit(1).get();
511 sem.acquire();
512 final Thread blocker = new Thread(() -> {
513 try {
514 sem.acquire();
515 } catch (final InterruptedException e) {
516 Thread.currentThread().interrupt();
517 } catch (final IllegalStateException e) {
518
519 }
520 }, "testShutdownWakesBlockedAcquireThreads");
521 blocker.setDaemon(true);
522 blocker.start();
523
524 final long parkDeadline = System.nanoTime() + Duration.ofSeconds(2).toNanos();
525 while (System.nanoTime() < parkDeadline && blocker.getState() != Thread.State.WAITING) {
526 Thread.sleep(10);
527 }
528 sem.shutdown();
529
530
531
532 blocker.join(5000);
533 assertFalse(blocker.isAlive(), "TimedSemaphore.shutdown() failed to wake thread blocked in acquire(): blocker still alive in state="
534 + blocker.getState() + " 5s after shutdown(). Bug present (shutdown() does not call notifyAll()).");
535 });
536 }
537
538
539
540
541
542
543 @Test
544 void testStartTimer() throws InterruptedException {
545 final TimedSemaphoreTestImpl semaphore = new TimedSemaphoreTestImpl(PERIOD_MILLIS, UNIT, LIMIT);
546 final ScheduledFuture<?> future = semaphore.startTimer();
547 assertNotNull(future, "No future returned");
548 ThreadUtils.sleepQuietly(DURATION);
549 final int trials = 10;
550 int count = 0;
551 do {
552 Thread.sleep(PERIOD_MILLIS);
553 assertFalse(count++ > trials, "endOfPeriod() not called!");
554 } while (semaphore.getPeriodEnds() <= 0);
555 semaphore.shutdown();
556 }
557
558
559
560
561
562 @Test
563 void testTryAcquire() throws InterruptedException {
564 final TimedSemaphore semaphore = new TimedSemaphore(PERIOD_MILLIS, TimeUnit.SECONDS, LIMIT);
565 final TryAcquireThread[] threads = new TryAcquireThread[3 * LIMIT];
566 final CountDownLatch latch = new CountDownLatch(1);
567 for (int i = 0; i < threads.length; i++) {
568 threads[i] = new TryAcquireThread(semaphore, latch);
569 threads[i].start();
570 }
571 latch.countDown();
572 int permits = 0;
573 for (final TryAcquireThread t : threads) {
574 t.join();
575 if (t.acquired) {
576 permits++;
577 }
578 }
579 assertEquals(LIMIT, permits, "Wrong number of permits granted");
580 }
581
582
583
584
585 @Test
586 void testTryAcquireAfterShutdown() {
587 final TimedSemaphore semaphore = new TimedSemaphore(PERIOD_MILLIS, UNIT, LIMIT);
588 semaphore.shutdown();
589 assertThrows(IllegalStateException.class, semaphore::tryAcquire);
590 }
591 }