|
21 | 21 | import org.jetbrains.annotations.Nullable; |
22 | 22 | import org.junit.jupiter.api.Test; |
23 | 23 |
|
| 24 | +import java.util.concurrent.CyclicBarrier; |
| 25 | +import java.util.concurrent.ExecutorService; |
| 26 | +import java.util.concurrent.Executors; |
| 27 | +import java.util.concurrent.Future; |
24 | 28 | import java.util.concurrent.Semaphore; |
| 29 | +import java.util.concurrent.TimeUnit; |
25 | 30 | import java.util.concurrent.atomic.AtomicInteger; |
26 | 31 | import java.util.concurrent.atomic.AtomicReference; |
27 | 32 |
|
28 | 33 | import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; |
29 | 34 | import static org.junit.jupiter.api.Assertions.assertEquals; |
| 35 | +import static org.junit.jupiter.api.Assertions.assertThrows; |
30 | 36 | import static org.junit.jupiter.api.Assertions.assertTrue; |
31 | 37 |
|
32 | 38 | public class ReleasePermitOnCompleteTest { |
@@ -59,4 +65,115 @@ public void onThrowable(Throwable t) { |
59 | 65 | assertEquals(1, completedCalls.get()); |
60 | 66 | assertEquals(1, throwableCalls.get()); |
61 | 67 | } |
| 68 | + |
| 69 | + @Test |
| 70 | + public void releasesPermitOnceWhenOnCompletedThrows() { |
| 71 | + Semaphore available = new Semaphore(Integer.MAX_VALUE); |
| 72 | + assertTrue(available.tryAcquire()); |
| 73 | + |
| 74 | + AtomicInteger throwableCalls = new AtomicInteger(); |
| 75 | + AsyncHandler<Object> handler = new AsyncCompletionHandler<Object>() { |
| 76 | + @Override |
| 77 | + public @Nullable Object onCompleted(@Nullable Response response) { |
| 78 | + throw new RuntimeException("boom"); |
| 79 | + } |
| 80 | + |
| 81 | + @Override |
| 82 | + public void onThrowable(Throwable t) { |
| 83 | + throwableCalls.incrementAndGet(); |
| 84 | + } |
| 85 | + }; |
| 86 | + AsyncHandler<Object> wrapped = ReleasePermitOnComplete.wrap(handler, available); |
| 87 | + |
| 88 | + Throwable thrown = assertThrows(Exception.class, wrapped::onCompleted); |
| 89 | + assertDoesNotThrow(() -> wrapped.onThrowable(thrown)); |
| 90 | + |
| 91 | + assertEquals(Integer.MAX_VALUE, available.availablePermits()); |
| 92 | + assertEquals(1, throwableCalls.get()); |
| 93 | + } |
| 94 | + |
| 95 | + @Test |
| 96 | + public void releasesSinglePermitOnNormalCompletion() { |
| 97 | + Semaphore available = new Semaphore(Integer.MAX_VALUE); |
| 98 | + assertTrue(available.tryAcquire()); |
| 99 | + |
| 100 | + AsyncHandler<Object> handler = new AsyncCompletionHandler<Object>() { |
| 101 | + @Override |
| 102 | + public @Nullable Object onCompleted(@Nullable Response response) { |
| 103 | + return null; |
| 104 | + } |
| 105 | + }; |
| 106 | + AsyncHandler<Object> wrapped = ReleasePermitOnComplete.wrap(handler, available); |
| 107 | + |
| 108 | + assertDoesNotThrow(() -> wrapped.onCompleted()); |
| 109 | + |
| 110 | + assertEquals(Integer.MAX_VALUE, available.availablePermits()); |
| 111 | + } |
| 112 | + |
| 113 | + @Test |
| 114 | + public void releasesSinglePermitOnThrowableOnly() { |
| 115 | + Semaphore available = new Semaphore(Integer.MAX_VALUE); |
| 116 | + assertTrue(available.tryAcquire()); |
| 117 | + |
| 118 | + AsyncHandler<Object> handler = new AsyncCompletionHandler<Object>() { |
| 119 | + @Override |
| 120 | + public @Nullable Object onCompleted(@Nullable Response response) { |
| 121 | + return null; |
| 122 | + } |
| 123 | + }; |
| 124 | + AsyncHandler<Object> wrapped = ReleasePermitOnComplete.wrap(handler, available); |
| 125 | + |
| 126 | + wrapped.onThrowable(new RuntimeException("failed")); |
| 127 | + |
| 128 | + assertEquals(Integer.MAX_VALUE, available.availablePermits()); |
| 129 | + } |
| 130 | + |
| 131 | + @Test |
| 132 | + public void releasesPermitOnceUnderConcurrentTerminalCallbacks() throws Exception { |
| 133 | + Semaphore available = new Semaphore(Integer.MAX_VALUE); |
| 134 | + assertTrue(available.tryAcquire()); |
| 135 | + |
| 136 | + AsyncHandler<Object> handler = new AsyncCompletionHandler<Object>() { |
| 137 | + @Override |
| 138 | + public @Nullable Object onCompleted(@Nullable Response response) { |
| 139 | + return null; |
| 140 | + } |
| 141 | + }; |
| 142 | + AsyncHandler<Object> wrapped = ReleasePermitOnComplete.wrap(handler, available); |
| 143 | + |
| 144 | + CyclicBarrier barrier = new CyclicBarrier(2); |
| 145 | + ExecutorService executor = Executors.newFixedThreadPool(2); |
| 146 | + try { |
| 147 | + Future<?> completion = executor.submit(() -> { |
| 148 | + awaitBarrier(barrier); |
| 149 | + completeQuietly(wrapped); |
| 150 | + }); |
| 151 | + Future<?> failure = executor.submit(() -> { |
| 152 | + awaitBarrier(barrier); |
| 153 | + wrapped.onThrowable(new RuntimeException("racing failure")); |
| 154 | + }); |
| 155 | + completion.get(5, TimeUnit.SECONDS); |
| 156 | + failure.get(5, TimeUnit.SECONDS); |
| 157 | + } finally { |
| 158 | + executor.shutdownNow(); |
| 159 | + } |
| 160 | + |
| 161 | + assertEquals(Integer.MAX_VALUE, available.availablePermits()); |
| 162 | + } |
| 163 | + |
| 164 | + private static void awaitBarrier(CyclicBarrier barrier) { |
| 165 | + try { |
| 166 | + barrier.await(5, TimeUnit.SECONDS); |
| 167 | + } catch (Exception e) { |
| 168 | + throw new RuntimeException(e); |
| 169 | + } |
| 170 | + } |
| 171 | + |
| 172 | + private static void completeQuietly(AsyncHandler<Object> handler) { |
| 173 | + try { |
| 174 | + handler.onCompleted(); |
| 175 | + } catch (Exception e) { |
| 176 | + throw new RuntimeException(e); |
| 177 | + } |
| 178 | + } |
62 | 179 | } |
0 commit comments