3636 * It manages the lifecycle of the original attempt and any subsequent hedged attempts.
3737 */
3838class CancellationSharer extends AbstractApiFuture <PublishResponse > {
39- private final Publisher .OutstandingBatch batch ;
39+ private Publisher .OutstandingBatch batch ;
4040 private final Publisher publisher ;
4141
4242 // Guarded by lock
@@ -47,6 +47,11 @@ class CancellationSharer extends AbstractApiFuture<PublishResponse> {
4747 private final Lock lock = new ReentrantLock ();
4848 private final AtomicBoolean isInQueue = new AtomicBoolean (false );
4949
50+ private void cleanupLocked () {
51+ runningAttempts .clear ();
52+ this .batch = null ;
53+ }
54+
5055 CancellationSharer (final Publisher .OutstandingBatch batch , final Publisher publisher ) {
5156 this .batch = batch ;
5257 this .publisher = publisher ;
@@ -90,19 +95,18 @@ private void handleAttemptSuccess(final int attemptNumber, final PublishResponse
9095 batch .successfulAttempt = attemptNumber ;
9196 set (response );
9297 cancelAllExceptLocked (attemptNumber );
98+ cleanupLocked ();
9399 } finally {
94100 lock .unlock ();
95101 }
96102 publisher .refillTokenBucket ();
97103 }
98104
99105 private void handleAttemptFailure (final int attemptNumber , final Throwable t ) {
100- boolean shouldRemoveFromQueue = false ;
101106 lock .lock ();
102107 try {
103108 if (done ) {
104- return ; // <-- Exit early before modifying runningAttempts to avoid
105- // ConcurrentModificationException
109+ return ;
106110 }
107111 runningAttempts .remove (attemptNumber );
108112 lastError = t ;
@@ -117,17 +121,11 @@ private void handleAttemptFailure(final int attemptNumber, final Throwable t) {
117121 done = true ;
118122 setException (lastError );
119123 cancelAllLocked ();
120- if (isInQueue .get ()) {
121- shouldRemoveFromQueue = true ;
122- }
124+ cleanupLocked ();
123125 }
124126 } finally {
125127 lock .unlock ();
126128 }
127-
128- if (shouldRemoveFromQueue ) {
129- publisher .removeFromHedgingQueue (this );
130- }
131129 }
132130
133131 void checkCompletionOnQueueExit () {
@@ -139,6 +137,7 @@ void checkCompletionOnQueueExit() {
139137 lastError != null
140138 ? lastError
141139 : new RuntimeException ("Hedging failed with no active attempts" ));
140+ cleanupLocked ();
142141 }
143142 } finally {
144143 lock .unlock ();
@@ -148,24 +147,17 @@ void checkCompletionOnQueueExit() {
148147 @ Override
149148 public boolean cancel (final boolean mayInterruptIfRunning ) {
150149 boolean cancelled = false ;
151- boolean shouldRemoveFromQueue = false ;
152150 lock .lock ();
153151 try {
154152 if (super .cancel (mayInterruptIfRunning )) {
155153 cancelled = true ;
156154 done = true ;
157- if (isInQueue .get ()) {
158- shouldRemoveFromQueue = true ;
159- }
160155 cancelAllLocked ();
156+ cleanupLocked ();
161157 }
162158 } finally {
163159 lock .unlock ();
164160 }
165-
166- if (shouldRemoveFromQueue ) {
167- publisher .removeFromHedgingQueue (this );
168- }
169161 return cancelled ;
170162 }
171163
@@ -190,7 +182,12 @@ AtomicBoolean isInQueue() {
190182 return isInQueue ;
191183 }
192184
193- Publisher .OutstandingBatch getBatch () {
194- return batch ;
185+ Publisher .OutstandingBatch getBatchIfActive () {
186+ lock .lock ();
187+ try {
188+ return done ? null : batch ;
189+ } finally {
190+ lock .unlock ();
191+ }
195192 }
196193}
0 commit comments