Skip to content

Commit ff1b8f2

Browse files
Implement releaseActiveOperationLock
1 parent 30a4292 commit ff1b8f2

2 files changed

Lines changed: 22 additions & 0 deletions

File tree

server/src/main/java/org/apache/cassandra/sidecar/job/OperationalJobManager.java

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -136,6 +136,7 @@ public void trySubmitJob(OperationalJob job,
136136
{
137137
if (ar.succeeded() && Boolean.TRUE.equals(ar.result()))
138138
{
139+
job.asyncResult().onComplete(result -> releaseActiveOperationLock(job));
139140
trackAndExecute(job, onComplete, serviceExecutorPool, waitTime);
140141
}
141142
else
@@ -207,6 +208,24 @@ private Future<Boolean> acquireActiveOperationLock(OperationalJob job)
207208
return internalExecutorPool.executeBlocking(() -> coordinator.trySetActive(job.operationType(), job.jobId()), false);
208209
}
209210

211+
/**
212+
* Releases the active operation lock previously acquired for the given job. The release performs
213+
* blocking storage I/O, so it runs on the internal executor pool off the event loop. A failure to clear is logged
214+
* rather than surfaced, since the job has already completed by this point.
215+
*
216+
* @param job the job whose active operation lock should be released
217+
*/
218+
private void releaseActiveOperationLock(OperationalJob job)
219+
{
220+
if (coordinator == null)
221+
{
222+
return;
223+
}
224+
internalExecutorPool.executeBlocking(() -> coordinator.clearActive(job.operationType(), job.jobId()), false)
225+
.onFailure(e -> logger.error("Failed to clear active operation lock. jobId={} operationType={}",
226+
job.jobId(), job.operationType(), e));
227+
}
228+
210229
/**
211230
* Builds the conflict exception describing why the active operation lock could not be acquired.
212231
*

server/src/test/java/org/apache/cassandra/sidecar/job/OperationalJobManagerTest.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,7 @@
4747
import static org.mockito.ArgumentMatchers.any;
4848
import static org.mockito.Mockito.mock;
4949
import static org.mockito.Mockito.never;
50+
import static org.mockito.Mockito.timeout;
5051
import static org.mockito.Mockito.verify;
5152
import static org.mockito.Mockito.when;
5253

@@ -199,6 +200,7 @@ void testCoordinatorCalledWhenJobRequiresCoordination() throws InterruptedExcept
199200
manager.trySubmitJob(job, onComplete, executorPool.service(), SecondBoundConfiguration.parse("5s"));
200201
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
201202
verify(coordinator).trySetActive(OperationType.MOVE, job.jobId());
203+
verify(coordinator, timeout(5000)).clearActive(OperationType.MOVE, job.jobId());
202204
}
203205

204206
@Test
@@ -224,6 +226,7 @@ void testConflictWhenCoordinatorReturnsFalse() throws InterruptedException
224226
assertThat(tracked.status()).isEqualTo(FAILED);
225227
assertThat(tracked.failureReason()).contains("An active operation already exists");
226228
assertThat(tracker.inflightJobsByOperation(job.name())).doesNotContain(job);
229+
verify(coordinator, never()).clearActive(any(), any());
227230
}
228231

229232
@Test

0 commit comments

Comments
 (0)