@@ -28,9 +28,15 @@ public class ExternalStorageConcurrencyTest {
2828 /** Takes a permit around every request and blocks, so peak concurrency is observable. */
2929 private static final class PermittingDriver implements StorageDriver {
3030 private final CompletableFuture <Void > gate = new CompletableFuture <>();
31+ private final int expectedOperations ;
32+ private final CompletableFuture <Void > expectedOperationsStarted = new CompletableFuture <>();
3133 private final AtomicInteger inFlight = new AtomicInteger ();
3234 private final AtomicInteger peak = new AtomicInteger ();
3335
36+ private PermittingDriver (int expectedOperations ) {
37+ this .expectedOperations = expectedOperations ;
38+ }
39+
3440 @ Nonnull
3541 @ Override
3642 public String getName () {
@@ -46,13 +52,20 @@ public String getType() {
4652 private <T > CompletableFuture <T > hold (T value ) {
4753 int current = inFlight .incrementAndGet ();
4854 peak .accumulateAndGet (current , Math ::max );
55+ if (current == expectedOperations ) {
56+ expectedOperationsStarted .complete (null );
57+ }
4958 return gate .thenApply (
5059 ignored -> {
5160 inFlight .decrementAndGet ();
5261 return value ;
5362 });
5463 }
5564
65+ private void awaitExpectedOperations () throws Exception {
66+ expectedOperationsStarted .get (5 , TimeUnit .SECONDS );
67+ }
68+
5669 @ Nonnull
5770 @ Override
5871 public CompletableFuture <List <StorageDriverClaim >> store (
@@ -108,23 +121,15 @@ private static List<Payload> payloads(int count) {
108121 return out ;
109122 }
110123
111- private static void awaitPeak (PermittingDriver driver , int expected ) throws Exception {
112- long deadline = System .nanoTime () + TimeUnit .SECONDS .toNanos (5 );
113- while (driver .peak .get () < expected && System .nanoTime () < deadline ) {
114- Thread .sleep (1 );
115- }
116- Thread .sleep (50 );
117- }
118-
119124 @ Test
120125 public void maxOperationsPerMessageBoundsOperations () throws Exception {
121- PermittingDriver driver = new PermittingDriver ();
126+ PermittingDriver driver = new PermittingDriver (3 );
122127 MessageStorageLimits limits = new MessageStorageLimits (3 , new AsyncSemaphore (100 ));
123128
124129 CompletableFuture <List <Payload >> result =
125130 transformer (driver ).store (payloads (6 ), null , CancellationToken .none (), limits );
126131
127- awaitPeak ( driver , 3 );
132+ driver . awaitExpectedOperations ( );
128133 assertEquals (3 , driver .peak .get ());
129134
130135 driver .gate .complete (null );
@@ -134,7 +139,7 @@ public void maxOperationsPerMessageBoundsOperations() throws Exception {
134139
135140 @ Test
136141 public void eachMessageGetsItsOwnBudget () throws Exception {
137- PermittingDriver driver = new PermittingDriver ();
142+ PermittingDriver driver = new PermittingDriver (4 );
138143 AsyncSemaphore shared = new AsyncSemaphore (100 );
139144 ExternalStoragePayloadTransformer transformer = transformer (driver );
140145
@@ -145,7 +150,7 @@ public void eachMessageGetsItsOwnBudget() throws Exception {
145150 transformer .store (
146151 payloads (3 ), null , CancellationToken .none (), new MessageStorageLimits (2 , shared )));
147152
148- awaitPeak ( driver , 4 );
153+ driver . awaitExpectedOperations ( );
149154 assertEquals ("two messages at 2 each, not 2 shared" , 4 , driver .peak .get ());
150155
151156 driver .gate .complete (null );
@@ -156,7 +161,7 @@ public void eachMessageGetsItsOwnBudget() throws Exception {
156161
157162 @ Test
158163 public void maxDriverOperationsSharedAcrossMessages () throws Exception {
159- PermittingDriver driver = new PermittingDriver ();
164+ PermittingDriver driver = new PermittingDriver (3 );
160165 AsyncSemaphore shared = new AsyncSemaphore (3 );
161166 ExternalStoragePayloadTransformer transformer = transformer (driver );
162167
@@ -167,7 +172,7 @@ public void maxDriverOperationsSharedAcrossMessages() throws Exception {
167172 transformer .store (
168173 payloads (4 ), null , CancellationToken .none (), new MessageStorageLimits (10 , shared )));
169174
170- awaitPeak ( driver , 3 );
175+ driver . awaitExpectedOperations ( );
171176 assertEquals ("one instance-wide budget spans both messages" , 3 , driver .peak .get ());
172177
173178 driver .gate .complete (null );
0 commit comments