diff --git a/graylog-storage-elasticsearch7/src/main/java/org/graylog/storage/elasticsearch7/ClusterAdapterES7.java b/graylog-storage-elasticsearch7/src/main/java/org/graylog/storage/elasticsearch7/ClusterAdapterES7.java index 634e68f36b9f..dd01e8c45813 100644 --- a/graylog-storage-elasticsearch7/src/main/java/org/graylog/storage/elasticsearch7/ClusterAdapterES7.java +++ b/graylog-storage-elasticsearch7/src/main/java/org/graylog/storage/elasticsearch7/ClusterAdapterES7.java @@ -327,6 +327,13 @@ public ShardStats shardStats() { .orElseThrow(() -> new ElasticsearchException("Unable to retrieve shard stats.")); } + @Override + public int countOfClusterManagerEligibleNodes() { + return (int) nodesInfo().values().stream() + .filter(node -> node.roles().contains("cluster_manager") || node.roles().contains("master")) + .count(); + } + private Optional clusterHealth() { try { final ClusterHealthRequest request = new ClusterHealthRequest() diff --git a/graylog-storage-opensearch2/src/main/java/org/graylog/storage/opensearch2/ClusterAdapterOS2.java b/graylog-storage-opensearch2/src/main/java/org/graylog/storage/opensearch2/ClusterAdapterOS2.java index 710da232dec6..a158fc83c0e5 100644 --- a/graylog-storage-opensearch2/src/main/java/org/graylog/storage/opensearch2/ClusterAdapterOS2.java +++ b/graylog-storage-opensearch2/src/main/java/org/graylog/storage/opensearch2/ClusterAdapterOS2.java @@ -336,6 +336,13 @@ public ShardStats shardStats() { .orElseThrow(() -> new ElasticsearchException("Unable to retrieve shard stats.")); } + @Override + public int countOfClusterManagerEligibleNodes() { + return (int)nodesInfo().values().stream() + .filter(node -> node.roles().contains("cluster_manager") || node.roles().contains("master")) + .count(); + } + private Optional clusterHealth() { try { final ClusterHealthRequest request = new ClusterHealthRequest() diff --git a/graylog-storage-opensearch3/src/main/java/org/graylog/storage/opensearch3/ClusterAdapterOS.java b/graylog-storage-opensearch3/src/main/java/org/graylog/storage/opensearch3/ClusterAdapterOS.java index 683f0ddc656d..763ecfaf9f5b 100644 --- a/graylog-storage-opensearch3/src/main/java/org/graylog/storage/opensearch3/ClusterAdapterOS.java +++ b/graylog-storage-opensearch3/src/main/java/org/graylog/storage/opensearch3/ClusterAdapterOS.java @@ -360,6 +360,13 @@ public ShardStats shardStats() { .orElseThrow(() -> new ElasticsearchException("Unable to retrieve shard stats.")); } + @Override + public int countOfClusterManagerEligibleNodes() { + return (int)nodesInfo().values().stream() + .filter(node -> node.roles().contains("cluster_manager") || node.roles().contains("master")) + .count(); + } + private Optional clusterHealth() { final Time timeout = new Time.Builder().time(requestTimeout.toSeconds() + "s").build(); try { diff --git a/graylog-storage-opensearch3/src/test/java/org/graylog/storage/opensearch3/ClusterAdapterOSTest.java b/graylog-storage-opensearch3/src/test/java/org/graylog/storage/opensearch3/ClusterAdapterOSTest.java index 007e404b32ad..64656753f02f 100644 --- a/graylog-storage-opensearch3/src/test/java/org/graylog/storage/opensearch3/ClusterAdapterOSTest.java +++ b/graylog-storage-opensearch3/src/test/java/org/graylog/storage/opensearch3/ClusterAdapterOSTest.java @@ -18,6 +18,7 @@ import com.github.joschi.jadconfig.util.Duration; import com.google.common.io.Resources; +import org.assertj.core.api.Assertions; import org.graylog.storage.opensearch3.testing.client.mock.ServerlessOpenSearchClient; import org.graylog2.indexer.cluster.health.ClusterShardAllocation; import org.graylog2.indexer.cluster.health.NodeDiskUsageStats; @@ -45,6 +46,7 @@ void setUp() { final OfficialOpensearchClient client = ServerlessOpenSearchClient.builder() .stubResponse("GET", "/_nodes/*", Resources.getResource("nodes-response-without-host-field.json")) + .stubResponse("GET", "/_nodes", Resources.getResource("nodes-response-without-host-field.json")) .stubResponse("GET", "/_cat/nodes", Resources.getResource("cat_nodes.json")) .stubResponse("GET", "/_cluster/settings", Resources.getResource("cluster_settings.json")) .stubResponse("GET", "/_cat/allocation", Resources.getResource("cat_allocation.json")) @@ -138,4 +140,10 @@ void testClusterShardAllocation() { .extracting(NodeShardAllocation::shards) .containsExactly(15, 16); } + + @Test + void testManagerEligibleNodesCount() { + Assertions.assertThat(clusterAdapter.countOfClusterManagerEligibleNodes()) + .isEqualTo(3); // there are 3 nodes with role "master" in nodes-response-without-host-field.json + } } diff --git a/graylog2-server/src/main/java/org/graylog2/indexer/cluster/ClusterAdapter.java b/graylog2-server/src/main/java/org/graylog2/indexer/cluster/ClusterAdapter.java index 11c923dcc7aa..8a7b42f9e0d9 100644 --- a/graylog2-server/src/main/java/org/graylog2/indexer/cluster/ClusterAdapter.java +++ b/graylog2-server/src/main/java/org/graylog2/indexer/cluster/ClusterAdapter.java @@ -66,5 +66,10 @@ public interface ClusterAdapter { ShardStats shardStats(); + /** + * The cluster health response has no such field, so implementations derive it from each node's roles. + */ + int countOfClusterManagerEligibleNodes(); + Optional deflectorHealth(Collection indices); }