|
21 | 21 | import org.elasticsearch.ResourceNotFoundException;
|
22 | 22 | import org.elasticsearch.Version;
|
23 | 23 | import org.elasticsearch.client.transport.TransportClient;
|
| 24 | +import org.elasticsearch.cluster.ClusterName; |
| 25 | +import org.elasticsearch.cluster.ClusterState; |
24 | 26 | import org.elasticsearch.cluster.Diff;
|
25 | 27 | import org.elasticsearch.cluster.NamedDiff;
|
26 | 28 | import org.elasticsearch.cluster.metadata.MetaData;
|
27 | 29 | import org.elasticsearch.cluster.metadata.MetaData.Custom;
|
| 30 | +import org.elasticsearch.cluster.node.DiscoveryNode; |
| 31 | +import org.elasticsearch.cluster.node.DiscoveryNodes; |
28 | 32 | import org.elasticsearch.common.ParseField;
|
29 | 33 | import org.elasticsearch.common.UUIDs;
|
30 | 34 | import org.elasticsearch.common.bytes.BytesReference;
|
|
33 | 37 | import org.elasticsearch.common.io.stream.NamedWriteableRegistry;
|
34 | 38 | import org.elasticsearch.common.io.stream.NamedWriteableRegistry.Entry;
|
35 | 39 | import org.elasticsearch.common.io.stream.StreamInput;
|
| 40 | +import org.elasticsearch.common.io.stream.StreamOutput; |
36 | 41 | import org.elasticsearch.common.io.stream.Writeable;
|
37 | 42 | import org.elasticsearch.common.xcontent.NamedXContentRegistry;
|
38 | 43 | import org.elasticsearch.common.xcontent.ToXContent;
|
| 44 | +import org.elasticsearch.common.xcontent.XContentBuilder; |
39 | 45 | import org.elasticsearch.common.xcontent.XContentFactory;
|
40 | 46 | import org.elasticsearch.common.xcontent.XContentParser;
|
41 | 47 | import org.elasticsearch.common.xcontent.XContentType;
|
|
65 | 71 | import static org.elasticsearch.test.VersionUtils.getPreviousVersion;
|
66 | 72 | import static org.elasticsearch.test.VersionUtils.randomVersionBetween;
|
67 | 73 | import static org.hamcrest.Matchers.equalTo;
|
| 74 | +import static org.hamcrest.Matchers.not; |
| 75 | +import static org.hamcrest.Matchers.sameInstance; |
68 | 76 |
|
69 | 77 | public class PersistentTasksCustomMetaDataTests extends AbstractDiffableSerializationTestCase<Custom> {
|
70 | 78 |
|
@@ -305,6 +313,91 @@ public void testFeatureSerialization() throws IOException {
|
305 | 313 | assertThat(read.taskMap().keySet(), equalTo(Collections.singleton("test_compatible")));
|
306 | 314 | }
|
307 | 315 |
|
| 316 | + public void testDisassociateDeadNodes_givenNoPersistentTasks() { |
| 317 | + ClusterState originalState = ClusterState.builder(new ClusterName("persistent-tasks-tests")).build(); |
| 318 | + ClusterState returnedState = PersistentTasksCustomMetaData.deassociateDeadNodes(originalState); |
| 319 | + assertThat(originalState, sameInstance(returnedState)); |
| 320 | + } |
| 321 | + |
| 322 | + public void testDisassociateDeadNodes_givenAssignedPersistentTask() { |
| 323 | + DiscoveryNodes nodes = DiscoveryNodes.builder() |
| 324 | + .add(new DiscoveryNode("node1", buildNewFakeTransportAddress(), Version.CURRENT)) |
| 325 | + .localNodeId("node1") |
| 326 | + .masterNodeId("node1") |
| 327 | + .build(); |
| 328 | + |
| 329 | + String taskName = "test/task"; |
| 330 | + PersistentTasksCustomMetaData.Builder tasksBuilder = PersistentTasksCustomMetaData.builder() |
| 331 | + .addTask("task-id", taskName, emptyTaskParams(taskName), |
| 332 | + new PersistentTasksCustomMetaData.Assignment("node1", "test assignment")); |
| 333 | + |
| 334 | + ClusterState originalState = ClusterState.builder(new ClusterName("persistent-tasks-tests")) |
| 335 | + .nodes(nodes) |
| 336 | + .metaData(MetaData.builder().putCustom(PersistentTasksCustomMetaData.TYPE, tasksBuilder.build())) |
| 337 | + .build(); |
| 338 | + ClusterState returnedState = PersistentTasksCustomMetaData.deassociateDeadNodes(originalState); |
| 339 | + assertThat(originalState, sameInstance(returnedState)); |
| 340 | + |
| 341 | + PersistentTasksCustomMetaData originalTasks = PersistentTasksCustomMetaData.getPersistentTasksCustomMetaData(originalState); |
| 342 | + PersistentTasksCustomMetaData returnedTasks = PersistentTasksCustomMetaData.getPersistentTasksCustomMetaData(returnedState); |
| 343 | + assertEquals(originalTasks, returnedTasks); |
| 344 | + } |
| 345 | + |
| 346 | + public void testDisassociateDeadNodes() { |
| 347 | + DiscoveryNodes nodes = DiscoveryNodes.builder() |
| 348 | + .add(new DiscoveryNode("node1", buildNewFakeTransportAddress(), Version.CURRENT)) |
| 349 | + .localNodeId("node1") |
| 350 | + .masterNodeId("node1") |
| 351 | + .build(); |
| 352 | + |
| 353 | + String taskName = "test/task"; |
| 354 | + PersistentTasksCustomMetaData.Builder tasksBuilder = PersistentTasksCustomMetaData.builder() |
| 355 | + .addTask("assigned-task", taskName, emptyTaskParams(taskName), |
| 356 | + new PersistentTasksCustomMetaData.Assignment("node1", "test assignment")) |
| 357 | + .addTask("task-on-deceased-node", taskName, emptyTaskParams(taskName), |
| 358 | + new PersistentTasksCustomMetaData.Assignment("left-the-cluster", "test assignment")); |
| 359 | + |
| 360 | + ClusterState originalState = ClusterState.builder(new ClusterName("persistent-tasks-tests")) |
| 361 | + .nodes(nodes) |
| 362 | + .metaData(MetaData.builder().putCustom(PersistentTasksCustomMetaData.TYPE, tasksBuilder.build())) |
| 363 | + .build(); |
| 364 | + ClusterState returnedState = PersistentTasksCustomMetaData.deassociateDeadNodes(originalState); |
| 365 | + assertThat(originalState, not(sameInstance(returnedState))); |
| 366 | + |
| 367 | + PersistentTasksCustomMetaData originalTasks = PersistentTasksCustomMetaData.getPersistentTasksCustomMetaData(originalState); |
| 368 | + PersistentTasksCustomMetaData returnedTasks = PersistentTasksCustomMetaData.getPersistentTasksCustomMetaData(returnedState); |
| 369 | + assertNotEquals(originalTasks, returnedTasks); |
| 370 | + |
| 371 | + assertEquals(originalTasks.getTask("assigned-task"), returnedTasks.getTask("assigned-task")); |
| 372 | + assertNotEquals(originalTasks.getTask("task-on-deceased-node"), returnedTasks.getTask("task-on-deceased-node")); |
| 373 | + assertEquals(PersistentTasksCustomMetaData.LOST_NODE_ASSIGNMENT, returnedTasks.getTask("task-on-deceased-node").getAssignment()); |
| 374 | + } |
| 375 | + |
| 376 | + private PersistentTaskParams emptyTaskParams(String taskName) { |
| 377 | + return new PersistentTaskParams() { |
| 378 | + |
| 379 | + @Override |
| 380 | + public XContentBuilder toXContent(XContentBuilder builder, Params params) { |
| 381 | + return builder; |
| 382 | + } |
| 383 | + |
| 384 | + @Override |
| 385 | + public void writeTo(StreamOutput out) { |
| 386 | + |
| 387 | + } |
| 388 | + |
| 389 | + @Override |
| 390 | + public String getWriteableName() { |
| 391 | + return taskName; |
| 392 | + } |
| 393 | + |
| 394 | + @Override |
| 395 | + public Version getMinimalSupportedVersion() { |
| 396 | + return Version.CURRENT; |
| 397 | + } |
| 398 | + }; |
| 399 | + } |
| 400 | + |
308 | 401 | private Assignment randomAssignment() {
|
309 | 402 | if (randomBoolean()) {
|
310 | 403 | if (randomBoolean()) {
|
|
0 commit comments