Skip to content

Commit 509b15b

Browse files
Implement DisabledOperationalJobCoordinator
1 parent ff1b8f2 commit 509b15b

6 files changed

Lines changed: 155 additions & 36 deletions

File tree

Lines changed: 68 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,68 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing, software
13+
* distributed under the License is distributed on an "AS IS" BASIS,
14+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15+
* See the License for the specific language governing permissions and
16+
* limitations under the License.
17+
*/
18+
19+
package org.apache.cassandra.sidecar.job;
20+
21+
import java.util.Collections;
22+
import java.util.Map;
23+
import java.util.UUID;
24+
25+
import org.apache.cassandra.sidecar.common.data.OperationType;
26+
import org.jetbrains.annotations.NotNull;
27+
import org.jetbrains.annotations.Nullable;
28+
29+
/**
30+
* An {@link OperationalJobCoordinator} used when a Sidecar instance is not configured to support
31+
* coordinated cluster-wide operations.
32+
* <p>
33+
* Jobs that do not require coordination ({@link OperationalJob#requiresCoordination()} returns
34+
* {@code false}) never reach this coordinator, so uncoordinated operations (e.g. decommission) are
35+
* unaffected.
36+
*/
37+
public class DisabledOperationalJobCoordinator implements OperationalJobCoordinator
38+
{
39+
private static final String NOT_SUPPORTED_MESSAGE =
40+
"Operational job coordination is not supported by this Sidecar instance. "
41+
+ "Configure a coordinator to enable coordinated cluster-wide operations.";
42+
43+
@Override
44+
public boolean trySetActive(OperationType operationType, UUID operationId)
45+
{
46+
throw new UnsupportedOperationException(NOT_SUPPORTED_MESSAGE);
47+
}
48+
49+
@Override
50+
public boolean clearActive(OperationType operationType, UUID operationId)
51+
{
52+
throw new UnsupportedOperationException(NOT_SUPPORTED_MESSAGE);
53+
}
54+
55+
@Override
56+
@Nullable
57+
public UUID getActiveOperation(OperationType operationType)
58+
{
59+
return null;
60+
}
61+
62+
@Override
63+
@NotNull
64+
public Map<OperationType, UUID> getActiveOperations()
65+
{
66+
return Collections.emptyMap();
67+
}
68+
}

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

Lines changed: 7 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -43,29 +43,20 @@ public class OperationalJobManager
4343
{
4444
protected final Logger logger = LoggerFactory.getLogger(this.getClass());
4545
private final OperationalJobTracker jobTracker;
46-
@Nullable
4746
private final OperationalJobCoordinator coordinator;
4847
private final TaskExecutorPool internalExecutorPool;
4948

5049
/**
51-
* Creates a manager instance without a coordinator.
52-
*
53-
* @param jobTracker the tracker for the operational jobs
54-
*/
55-
public OperationalJobManager(OperationalJobTracker jobTracker, ExecutorPools executorPools)
56-
{
57-
this(jobTracker, null, executorPools);
58-
}
59-
60-
/**
61-
* Creates a manager instance with a coordinator for cluster-wide operation mutual exclusion.
50+
* Creates a manager instance with a coordinator for cluster-wide operation mutual exclusion. Instances that
51+
* do not support coordination bind a {@link DisabledOperationalJobCoordinator}, which fails coordination
52+
* requests.
6253
*
6354
* @param jobTracker the tracker for the operational jobs
64-
* @param coordinator the coordinator for cluster-wide operations, or {@code null} if not needed
55+
* @param coordinator the coordinator for cluster-wide operations
6556
*/
6657
@Inject
6758
public OperationalJobManager(OperationalJobTracker jobTracker,
68-
@Nullable OperationalJobCoordinator coordinator,
59+
OperationalJobCoordinator coordinator,
6960
ExecutorPools executorPools)
7061
{
7162
this.jobTracker = jobTracker;
@@ -196,15 +187,11 @@ private void checkConflict(OperationalJob job) throws OperationalJobConflictExce
196187
*
197188
* @param job the job requiring coordination
198189
* @return a future resolving to {@code true} if the lock was acquired, {@code false} if another operation
199-
* already holds it, or a failed future if coordination could not be attempted (e.g. no coordinator
200-
* is configured or the storage call failed)
190+
* already holds it, or a failed future if coordination could not be attempted (e.g. coordination is
191+
* disabled on this instance or the storage call failed)
201192
*/
202193
private Future<Boolean> acquireActiveOperationLock(OperationalJob job)
203194
{
204-
if (coordinator == null)
205-
{
206-
return Future.failedFuture("Job requires coordination but no OperationalJobCoordinator is configured");
207-
}
208195
return internalExecutorPool.executeBlocking(() -> coordinator.trySetActive(job.operationType(), job.jobId()), false);
209196
}
210197

@@ -217,10 +204,6 @@ private Future<Boolean> acquireActiveOperationLock(OperationalJob job)
217204
*/
218205
private void releaseActiveOperationLock(OperationalJob job)
219206
{
220-
if (coordinator == null)
221-
{
222-
return;
223-
}
224207
internalExecutorPool.executeBlocking(() -> coordinator.clearActive(job.operationType(), job.jobId()), false)
225208
.onFailure(e -> logger.error("Failed to clear active operation lock. jobId={} operationType={}",
226209
job.jobId(), job.operationType(), e));

server/src/main/java/org/apache/cassandra/sidecar/modules/CassandraOperationsModule.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -61,10 +61,10 @@
6161
import org.apache.cassandra.sidecar.handlers.cassandra.NodeSettingsHandler;
6262
import org.apache.cassandra.sidecar.handlers.v2.cassandra.V2NodeSettingsHandler;
6363
import org.apache.cassandra.sidecar.handlers.validations.ValidateTableExistenceHandler;
64+
import org.apache.cassandra.sidecar.job.DisabledOperationalJobCoordinator;
6465
import org.apache.cassandra.sidecar.job.InMemoryOperationalJobTracker;
6566
import org.apache.cassandra.sidecar.job.OperationalJobCoordinator;
6667
import org.apache.cassandra.sidecar.job.OperationalJobTracker;
67-
import org.apache.cassandra.sidecar.job.StorageBackedOperationalJobCoordinator;
6868
import org.apache.cassandra.sidecar.modules.multibindings.KeyClassMapKey;
6969
import org.apache.cassandra.sidecar.modules.multibindings.TableSchemaMapKeys;
7070
import org.apache.cassandra.sidecar.modules.multibindings.VertxRouteMapKeys;
@@ -85,7 +85,7 @@ public class CassandraOperationsModule extends AbstractModule
8585
protected void configure()
8686
{
8787
bind(OperationalJobTracker.class).to(InMemoryOperationalJobTracker.class);
88-
bind(OperationalJobCoordinator.class).to(StorageBackedOperationalJobCoordinator.class);
88+
bind(OperationalJobCoordinator.class).to(DisabledOperationalJobCoordinator.class);
8989
}
9090

9191
@ProvidesIntoMap
Lines changed: 68 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,68 @@
1+
/*
2+
* Licensed to the Apache Software Foundation (ASF) under one
3+
* or more contributor license agreements. See the NOTICE file
4+
* distributed with this work for additional information
5+
* regarding copyright ownership. The ASF licenses this file
6+
* to you under the Apache License, Version 2.0 (the
7+
* "License"); you may not use this file except in compliance
8+
* with the License. You may obtain a copy of the License at
9+
*
10+
* http://www.apache.org/licenses/LICENSE-2.0
11+
*
12+
* Unless required by applicable law or agreed to in writing, software
13+
* distributed under the License is distributed on an "AS IS" BASIS,
14+
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
15+
* See the License for the specific language governing permissions and
16+
* limitations under the License.
17+
*/
18+
19+
package org.apache.cassandra.sidecar.job;
20+
21+
import java.util.UUID;
22+
23+
import org.junit.jupiter.api.Test;
24+
25+
import com.datastax.driver.core.utils.UUIDs;
26+
import org.apache.cassandra.sidecar.common.data.OperationType;
27+
28+
import static org.assertj.core.api.Assertions.assertThat;
29+
import static org.assertj.core.api.Assertions.assertThatThrownBy;
30+
31+
/**
32+
* Tests for {@link DisabledOperationalJobCoordinator}, the coordinator used when a Sidecar instance
33+
* does not support coordinated cluster-wide operations.
34+
*/
35+
class DisabledOperationalJobCoordinatorTest
36+
{
37+
private final OperationalJobCoordinator coordinator = new DisabledOperationalJobCoordinator();
38+
39+
@Test
40+
void testTrySetActiveThrows()
41+
{
42+
UUID operationId = UUIDs.timeBased();
43+
assertThatThrownBy(() -> coordinator.trySetActive(OperationType.MOVE, operationId))
44+
.isInstanceOf(UnsupportedOperationException.class)
45+
.hasMessageContaining("coordination is not supported by this Sidecar instance");
46+
}
47+
48+
@Test
49+
void testClearActiveThrows()
50+
{
51+
UUID operationId = UUIDs.timeBased();
52+
assertThatThrownBy(() -> coordinator.clearActive(OperationType.MOVE, operationId))
53+
.isInstanceOf(UnsupportedOperationException.class)
54+
.hasMessageContaining("coordination is not supported by this Sidecar instance");
55+
}
56+
57+
@Test
58+
void testGetActiveOperationReturnsNull()
59+
{
60+
assertThat(coordinator.getActiveOperation(OperationType.MOVE)).isNull();
61+
}
62+
63+
@Test
64+
void testGetActiveOperationsReturnsEmptyMap()
65+
{
66+
assertThat(coordinator.getActiveOperations()).isEmpty();
67+
}
68+
}

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

Lines changed: 9 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -82,7 +82,7 @@ void cleanup()
8282
void testWithNoDownstreamJob() throws InterruptedException
8383
{
8484
OperationalJobTracker tracker = new InMemoryOperationalJobTracker(4);
85-
OperationalJobManager manager = new OperationalJobManager(tracker, executorPool);
85+
OperationalJobManager manager = new OperationalJobManager(tracker, new DisabledOperationalJobCoordinator(), executorPool);
8686
CountDownLatch latch = new CountDownLatch(1);
8787

8888
OperationalJob testJob = OperationalJobTest.createOperationalJob(SUCCEEDED);
@@ -103,7 +103,7 @@ void testWithRunningDownstreamJob() throws InterruptedException
103103
{
104104
OperationalJob runningJob = OperationalJobTest.createOperationalJob(RUNNING);
105105
OperationalJobTracker tracker = new InMemoryOperationalJobTracker(4);
106-
OperationalJobManager manager = new OperationalJobManager(tracker, executorPool);
106+
OperationalJobManager manager = new OperationalJobManager(tracker, new DisabledOperationalJobCoordinator(), executorPool);
107107
CountDownLatch latch = new CountDownLatch(1);
108108

109109
BiConsumer<OperationalJob, OperationalJobConflictException> onComplete = (job, exception) -> {
@@ -121,7 +121,7 @@ void testWithLongRunningJob() throws InterruptedException
121121
{
122122
UUID jobId = UUIDs.timeBased();
123123
OperationalJobTracker tracker = new InMemoryOperationalJobTracker(4);
124-
OperationalJobManager manager = new OperationalJobManager(tracker, executorPool);
124+
OperationalJobManager manager = new OperationalJobManager(tracker, new DisabledOperationalJobCoordinator(), executorPool);
125125
CountDownLatch latch = new CountDownLatch(1);
126126

127127
OperationalJob testJob = OperationalJobTest.createOperationalJob(jobId, SecondBoundConfiguration.parse("2s"));
@@ -145,7 +145,7 @@ void testWithFailingJob() throws InterruptedException
145145
{
146146
UUID jobId = UUIDs.timeBased();
147147
OperationalJobTracker tracker = new InMemoryOperationalJobTracker(4);
148-
OperationalJobManager manager = new OperationalJobManager(tracker, executorPool);
148+
OperationalJobManager manager = new OperationalJobManager(tracker, new DisabledOperationalJobCoordinator(), executorPool);
149149
CountDownLatch latch = new CountDownLatch(1);
150150

151151
String msg = "Test Job failed";
@@ -230,17 +230,17 @@ void testConflictWhenCoordinatorReturnsFalse() throws InterruptedException
230230
}
231231

232232
@Test
233-
void testCoordinationFailsWhenNoCoordinatorConfigured() throws InterruptedException
233+
void testCoordinationFailsWhenCoordinationDisabled() throws InterruptedException
234234
{
235235
OperationalJobTracker tracker = new InMemoryOperationalJobTracker(4);
236-
// No coordinator is wired, yet the job requires coordination.
237-
OperationalJobManager manager = new OperationalJobManager(tracker, executorPool);
236+
// Coordination is disabled on this instance, yet the job requires coordination.
237+
OperationalJobManager manager = new OperationalJobManager(tracker, new DisabledOperationalJobCoordinator(), executorPool);
238238
CountDownLatch latch = new CountDownLatch(1);
239239

240240
OperationalJob job = createCoordinatedJob(UUIDs.timeBased());
241241
BiConsumer<OperationalJob, OperationalJobConflictException> onComplete = (j, ex) -> {
242242
assertThat(ex).isInstanceOf(OperationalJobConflictException.class);
243-
assertThat(ex.getMessage()).contains("no OperationalJobCoordinator is configured");
243+
assertThat(ex.getMessage()).contains("coordination is not supported by this Sidecar instance");
244244
latch.countDown();
245245
};
246246

@@ -249,7 +249,7 @@ void testCoordinationFailsWhenNoCoordinatorConfigured() throws InterruptedExcept
249249
OperationalJobInfo tracked = tracker.get(job.jobId());
250250
assertThat(tracked).isNotNull();
251251
assertThat(tracked.status()).isEqualTo(FAILED);
252-
assertThat(tracked.failureReason()).contains("no OperationalJobCoordinator is configured");
252+
assertThat(tracked.failureReason()).contains("coordination is not supported by this Sidecar instance");
253253
assertThat(tracker.inflightJobsByOperation(job.name())).doesNotContain(job);
254254
}
255255

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -234,7 +234,7 @@ void testMultipleRepairJobsRunningInParallel() throws Exception
234234
{
235235
// Create a job tracker and manager
236236
OperationalJobTracker tracker = new InMemoryOperationalJobTracker(10);
237-
OperationalJobManager manager = new OperationalJobManager(tracker, executorPool);
237+
OperationalJobManager manager = new OperationalJobManager(tracker, new DisabledOperationalJobCoordinator(), executorPool);
238238

239239
// Mock the storage operations
240240
StorageOperations storageOperations = mock(StorageOperations.class);

0 commit comments

Comments
 (0)