diff --git a/hadoop-hdds/docs/content/feature/Reconfigurability.md b/hadoop-hdds/docs/content/feature/Reconfigurability.md index 72537c901421..246a55227270 100644 --- a/hadoop-hdds/docs/content/feature/Reconfigurability.md +++ b/hadoop-hdds/docs/content/feature/Reconfigurability.md @@ -72,6 +72,8 @@ ozone admin reconfig --service=[OM|SCM|DATANODE] --address=` | - | Comma-separated HA SCM node IDs.| +| `ozone.scm.address..` | - | SCM RPC address in HA mode.| ### Storage Container Manager (SCM) diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java index 05bcbf6bffa2..cf85c4b5b184 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/proxy/SCMFailoverProxyProviderBase.java @@ -68,7 +68,7 @@ public abstract class SCMFailoverProxyProviderBase implements FailoverProxyPr // scmNodeId -> ProxyInfo private final Map> scmProxies; // scmNodeId -> SCMProxyInfo - private final Map scmProxyInfoMap; + private Map scmProxyInfoMap; private List scmNodeIds; // As SCM Client is shared across threads, performFailOver() @@ -122,7 +122,6 @@ public SCMFailoverProxyProviderBase(Class protocol, ConfigurationSource conf, this.scmVersion = RPC.getProtocolVersion(protocol); this.scmProxies = new HashMap<>(); - this.scmProxyInfoMap = new HashMap<>(); loadConfigs(); this.currentProxyIndex = 0; @@ -174,9 +173,19 @@ synchronized void replaceProxyInfoForTest(String nodeId, SCMProxyInfo info) { @VisibleForTesting protected synchronized void loadConfigs() { - List scmNodeInfoList = SCMNodeInfo.buildNodeInfo(conf); - scmNodeIds = new ArrayList<>(); + ScmProxyConfig newConfig = buildConfigs(); + scmNodeIds = newConfig.nodeIds; + scmProxyInfoMap = newConfig.proxyInfoMap; + } + /** + * Thread-safely resolves node addresses into an isolated holder. + * Fails atomically if any address is missing. + */ + private ScmProxyConfig buildConfigs() { + List scmNodeInfoList = SCMNodeInfo.buildNodeInfo(conf); + List newScmNodeIds = new ArrayList<>(); + Map newScmProxyInfoMap = new HashMap<>(); for (SCMNodeInfo scmNodeInfo : scmNodeInfoList) { String protocolAddress = getProtocolAddress(scmNodeInfo); @@ -188,15 +197,85 @@ protected synchronized void loadConfigs() { String scmServiceId = scmNodeInfo.getServiceId(); String scmNodeId = scmNodeInfo.getNodeId(); - scmNodeIds.add(scmNodeId); + newScmNodeIds.add(scmNodeId); // Preserve the original config string so DNS can be re-resolved // on connection failure when the SCM peer is rescheduled to a // new IP (Kubernetes pod-IP-change recovery). See // refreshProxyAddressIfChanged(String). SCMProxyInfo scmProxyInfo = new SCMProxyInfo(scmServiceId, scmNodeId, protocolAddr, protocolAddress); - scmProxyInfoMap.put(scmNodeId, scmProxyInfo); + newScmProxyInfoMap.put(scmNodeId, scmProxyInfo); + } + } + + return new ScmProxyConfig(newScmNodeIds, newScmProxyInfoMap); + } + + /** + * Reloads SCM nodes and proxies from updated config without a restart. + * Stops proxies for removed/changed nodes; fails atomically if config is invalid. + * In-flight calls on removed nodes persist until the next failover. + */ + public void changeConfig() { + // Resolve DNS before locking to prevent slow lookup from blocking callers. + ScmProxyConfig newConfig = buildConfigs(); + + Map> staleProxies = new HashMap<>(); + synchronized (this) { + Map oldProxyInfoMap = scmProxyInfoMap; + scmNodeIds = newConfig.nodeIds; + scmProxyInfoMap = newConfig.proxyInfoMap; + + // Re-sync proxy index to the new list, or fall back to first node if removed. + int newProxyIndex = scmNodeIds.indexOf(currentProxySCMNodeId); + if (newProxyIndex < 0) { + newProxyIndex = 0; + currentProxySCMNodeId = scmNodeIds.get(newProxyIndex); } + currentProxyIndex = newProxyIndex; + + // Drop removed failover target to prevent NPE on next failover. + if (updatedLeaderNodeID != null + && !scmProxyInfoMap.containsKey(updatedLeaderNodeID)) { + updatedLeaderNodeID = null; + } + + // Evict stale proxies under lock, but defer stopProxy until unlocked to avoid blocking. + for (Map.Entry entry : oldProxyInfoMap.entrySet()) { + String nodeId = entry.getKey(); + SCMProxyInfo newInfo = scmProxyInfoMap.get(nodeId); + if (newInfo == null + || !newInfo.getAddress().equals(entry.getValue().getAddress())) { + ProxyInfo staleProxy = scmProxies.remove(nodeId); + if (staleProxy != null && staleProxy.proxy != null) { + staleProxies.put(nodeId, staleProxy); + } + } + } + + getLogger().info("Reloaded SCM proxy configuration for protocol {} with {} nodes: {}", + protocolClass.getSimpleName(), scmNodeIds.size(), scmProxyInfoMap.values()); + } + + for (Map.Entry> entry : staleProxies.entrySet()) { + try { + RPC.stopProxy(entry.getValue().proxy); + } catch (RuntimeException stopEx) { + getLogger().warn("Failed to stop stale proxy for SCM node {}", + entry.getKey(), stopEx); + } + } + } + + /** Parsed node list and resolved addresses, built without touching shared state. */ + private static final class ScmProxyConfig { + private final List nodeIds; + private final Map proxyInfoMap; + + ScmProxyConfig(List nodeIds, + Map proxyInfoMap) { + this.nodeIds = nodeIds; + this.proxyInfoMap = proxyInfoMap; } } diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/HAUtils.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/HAUtils.java index 3ae1b451e1fc..403416981aa1 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/HAUtils.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/utils/HAUtils.java @@ -131,9 +131,19 @@ public static boolean addSCM(OzoneConfiguration conf, AddSCMRequest request, */ public static ScmBlockLocationProtocol getScmBlockClient( OzoneConfiguration conf) { + return getScmBlockClient(conf, + new SCMBlockLocationFailoverProxyProvider(conf)); + } + + /** + * Creates an SCM block client using the provided proxy provider. + * Retains the provider reference to support dynamic SCM node updates. + */ + public static ScmBlockLocationProtocol getScmBlockClient( + OzoneConfiguration conf, + SCMBlockLocationFailoverProxyProvider proxyProvider) { ScmBlockLocationProtocolClientSideTranslatorPB scmBlockLocationClient = - new ScmBlockLocationProtocolClientSideTranslatorPB( - new SCMBlockLocationFailoverProxyProvider(conf), conf); + new ScmBlockLocationProtocolClientSideTranslatorPB(proxyProvider, conf); return TracingUtil .createProxy(scmBlockLocationClient, ScmBlockLocationProtocol.class, conf); @@ -141,8 +151,17 @@ public static ScmBlockLocationProtocol getScmBlockClient( public static StorageContainerLocationProtocol getScmContainerClient( ConfigurationSource conf) { - SCMContainerLocationFailoverProxyProvider proxyProvider = - new SCMContainerLocationFailoverProxyProvider(conf, null); + return getScmContainerClient(conf, + new SCMContainerLocationFailoverProxyProvider(conf, null)); + } + + /** + * Creates an SCM container client using the provided proxy provider. + * Retains the provider reference to support dynamic SCM node updates. + */ + public static StorageContainerLocationProtocol getScmContainerClient( + ConfigurationSource conf, + SCMContainerLocationFailoverProxyProvider proxyProvider) { StorageContainerLocationProtocol scmContainerClient = TracingUtil.createProxy( new StorageContainerLocationProtocolClientSideTranslatorPB( diff --git a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/scm/proxy/TestSCMFailoverProxyProviderChangeConfig.java b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/scm/proxy/TestSCMFailoverProxyProviderChangeConfig.java new file mode 100644 index 000000000000..adc9ec154889 --- /dev/null +++ b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/scm/proxy/TestSCMFailoverProxyProviderChangeConfig.java @@ -0,0 +1,149 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hadoop.hdds.scm.proxy; + +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_ADDRESS_KEY; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_NODES_KEY; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_SERVICE_IDS_KEY; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNotSame; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.List; +import org.apache.hadoop.hdds.conf.ConfigurationException; +import org.apache.hadoop.hdds.conf.OzoneConfiguration; +import org.apache.hadoop.ozone.ha.ConfUtils; +import org.junit.jupiter.api.Test; + +/** + * Verifies that {@link SCMFailoverProxyProviderBase#changeConfig()} reloads the + * SCM node list from an updated configuration: adding and removing nodes, keeping + * the current proxy pointer valid, and leaving prior state intact if the new + * configuration is incomplete. + */ +public class TestSCMFailoverProxyProviderChangeConfig { + + private static final String SERVICE_ID = "scmservice"; + + private static OzoneConfiguration haConf(String nodes) { + OzoneConfiguration conf = new OzoneConfiguration(); + conf.set(OZONE_SCM_SERVICE_IDS_KEY, SERVICE_ID); + conf.set(ConfUtils.addSuffix(OZONE_SCM_NODES_KEY, SERVICE_ID), nodes); + return conf; + } + + private static void setAddress(OzoneConfiguration conf, String nodeId, + String host) { + conf.set(ConfUtils.addKeySuffixes(OZONE_SCM_ADDRESS_KEY, SERVICE_ID, nodeId), + host); + } + + @Test + public void testChangeConfigAddsNode() { + OzoneConfiguration conf = haConf("scm1,scm2"); + setAddress(conf, "scm1", "host1"); + setAddress(conf, "scm2", "host2"); + + SCMBlockLocationFailoverProxyProvider provider = + new SCMBlockLocationFailoverProxyProvider(conf); + assertEquals(2, provider.getSCMNodeIds().size()); + + // Operator adds a third SCM: its address key first, then the node list. + setAddress(conf, "scm3", "host3"); + conf.set(ConfUtils.addSuffix(OZONE_SCM_NODES_KEY, SERVICE_ID), + "scm1,scm2,scm3"); + provider.changeConfig(); + + List nodeIds = provider.getSCMNodeIds(); + assertEquals(3, nodeIds.size()); + assertTrue(nodeIds.contains("scm3")); + } + + @Test + public void testChangeConfigRemovesNode() { + OzoneConfiguration conf = haConf("scm1,scm2,scm3"); + setAddress(conf, "scm1", "host1"); + setAddress(conf, "scm2", "host2"); + setAddress(conf, "scm3", "host3"); + + SCMBlockLocationFailoverProxyProvider provider = + new SCMBlockLocationFailoverProxyProvider(conf); + + // Pass the preceding node so changeCurrentProxy advances to scm3 (to be removed). + List before = provider.getSCMNodeIds(); + int size = before.size(); + int scm3Index = before.indexOf("scm3"); + provider.changeCurrentProxy(before.get((scm3Index - 1 + size) % size)); + assertEquals("scm3", provider.getCurrentProxySCMNodeId()); + + conf.set(ConfUtils.addSuffix(OZONE_SCM_NODES_KEY, SERVICE_ID), "scm1,scm2"); + provider.changeConfig(); + + List nodeIds = provider.getSCMNodeIds(); + assertEquals(2, nodeIds.size()); + assertTrue(nodeIds.contains("scm1")); + assertTrue(nodeIds.contains("scm2")); + // The current proxy pointer must fall back to a still-configured node. + assertTrue(nodeIds.contains(provider.getCurrentProxySCMNodeId())); + } + + @Test + public void testAddressChangeEvictsCachedProxy() { + OzoneConfiguration conf = haConf("scm1,scm2"); + setAddress(conf, "scm1", "127.0.0.1"); + setAddress(conf, "scm2", "127.0.0.2"); + + SCMBlockLocationFailoverProxyProvider provider = + new SCMBlockLocationFailoverProxyProvider(conf); + + String current = provider.getCurrentProxySCMNodeId(); + Object first = provider.getProxy(); + // With no configuration change the cached proxy is reused. + assertSame(first, provider.getProxy()); + + // Change the current node's address: changeConfig must evict the stale proxy + // so the next getProxy() rebuilds against the new endpoint. + setAddress(conf, current, "127.0.0.3"); + provider.changeConfig(); + + assertNotSame(first, provider.getProxy()); + } + + @Test + public void testChangeConfigFailsWhenAddressMissing() { + OzoneConfiguration conf = haConf("scm1,scm2"); + setAddress(conf, "scm1", "host1"); + setAddress(conf, "scm2", "host2"); + + SCMBlockLocationFailoverProxyProvider provider = + new SCMBlockLocationFailoverProxyProvider(conf); + + // Node list references scm3 but its address is not set yet, so the reload + // must fail and leave the previous node set intact for a retry. + conf.set(ConfUtils.addSuffix(OZONE_SCM_NODES_KEY, SERVICE_ID), + "scm1,scm2,scm3"); + assertThrows(ConfigurationException.class, provider::changeConfig); + + List nodeIds = provider.getSCMNodeIds(); + assertEquals(2, nodeIds.size()); + assertTrue(nodeIds.contains("scm1")); + assertTrue(nodeIds.contains("scm2")); + } +} diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestOmSCMNodesReconfiguration.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestOmSCMNodesReconfiguration.java new file mode 100644 index 000000000000..7827e0986754 --- /dev/null +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestOmSCMNodesReconfiguration.java @@ -0,0 +1,363 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hadoop.hdds.scm; + +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_ADDRESS_KEY; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_NODES_KEY; +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import org.apache.hadoop.conf.ReconfigurationException; +import org.apache.hadoop.hdds.conf.OzoneConfiguration; +import org.apache.hadoop.hdds.conf.ReconfigurationHandler; +import org.apache.hadoop.hdds.scm.proxy.SCMFailoverProxyProviderBase; +import org.apache.hadoop.hdds.scm.proxy.SCMProxyInfo; +import org.apache.hadoop.hdds.scm.server.StorageContainerManager; +import org.apache.hadoop.ozone.MiniOzoneCluster; +import org.apache.hadoop.ozone.MiniOzoneHAClusterImpl; +import org.apache.hadoop.ozone.ha.ConfUtils; +import org.apache.hadoop.ozone.om.OzoneManager; +import org.apache.hadoop.ozone.om.ScmClient; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + +/** + * Tests the OM's SCM nodes reconfiguration wiring: the SCM node list and the + * per-node SCM addresses must be reconfigurable on a running OM, reloading its + * SCM failover proxies without a restart. Proxy-level add/remove behavior is + * covered by {@link org.apache.hadoop.hdds.scm.proxy.SCMFailoverProxyProviderBase} + * unit tests; this verifies the OM-side registration and callbacks end to end. + */ +@Timeout(300) +public class TestOmSCMNodesReconfiguration { + + private MiniOzoneHAClusterImpl cluster = null; + private String scmServiceId; + + @BeforeEach + public void init() throws Exception { + OzoneConfiguration conf = new OzoneConfiguration(); + scmServiceId = "scm-service-test1"; + cluster = MiniOzoneCluster.newHABuilder(conf) + .setOMServiceId("om-service-test1") + .setSCMServiceId(scmServiceId) + .setNumOfStorageContainerManagers(3) + .setNumOfOzoneManagers(1) + .build(); + cluster.waitForClusterToBeReady(); + } + + @AfterEach + public void shutdown() { + if (cluster != null) { + cluster.shutdown(); + } + } + + /** + * Ensures SCM node list and address prefix configurations are reconfigurable on the OM. + */ + @Test + void testScmNodesAndAddressReconfigurableOnOm() throws Exception { + ReconfigurationHandler handler = + cluster.getOzoneManager().getReconfigurationHandler(); + String scmNodesKey = + ConfUtils.addKeySuffixes(OZONE_SCM_NODES_KEY, scmServiceId); + + assertTrue(handler.isPropertyReconfigurable(scmNodesKey)); + assertTrue(handler.listReconfigureProperties().contains(scmNodesKey)); + + // Register address prefix so dynamically added nodes can be reconfigured. + for (StorageContainerManager scm : cluster.getStorageContainerManagers()) { + String scmAddrKey = ConfUtils.addKeySuffixes( + OZONE_SCM_ADDRESS_KEY, scmServiceId, scm.getSCMNodeId()); + assertTrue(handler.isPropertyReconfigurable(scmAddrKey)); + } + } + + /** + * Verifies that an empty SCM node list is rejected without modifying existing proxies. + */ + @Test + void testReconfigureScmNodesToBlankThrows() { + ReconfigurationHandler handler = + cluster.getOzoneManager().getReconfigurationHandler(); + String scmNodesKey = + ConfUtils.addKeySuffixes(OZONE_SCM_NODES_KEY, scmServiceId); + + assertThrows(ReconfigurationException.class, + () -> handler.reconfigureProperty(scmNodesKey, "")); + } + + /** + * Verifies that reconfiguring SCM nodes dynamically updates block and container failover proxies. + */ + @Test + void testReconfigureScmNodesReloadsProxies() throws Exception { + OzoneManager om = cluster.getOzoneManager(); + ReconfigurationHandler handler = om.getReconfigurationHandler(); + ScmClient scmClient = om.getScmClient(); + String scmNodesKey = + ConfUtils.addKeySuffixes(OZONE_SCM_NODES_KEY, scmServiceId); + + List before = + new ArrayList<>(scmClient.getContainerProxyProvider().getSCMNodeIds()); + assertEquals(3, before.size()); + + // Remove one SCM; its lingering address configuration allows reloading the remaining nodes. + String dropped = before.get(before.size() - 1); + List remaining = new ArrayList<>(before.subList(0, before.size() - 1)); + Set expected = new HashSet<>(remaining); + + handler.reconfigureProperty(scmNodesKey, String.join(",", remaining)); + + Set afterContainer = + new HashSet<>(scmClient.getContainerProxyProvider().getSCMNodeIds()); + Set afterBlock = + new HashSet<>(scmClient.getBlockProxyProvider().getSCMNodeIds()); + assertEquals(expected, afterContainer); + assertEquals(expected, afterBlock); + assertFalse(afterContainer.contains(dropped)); + + // Perform live RPCs to verify rebuilt proxies rather than just checking node IDs. + assertNotNull(scmClient.getContainerClient().getScmInfo()); + assertNotNull(scmClient.getBlockClient().getScmInfo()); + } + + /** + * Verifies that updating a per-node SCM address reloads failover proxies via the completion callback. + */ + @Test + void testReconfigureScmAddressReloadsProxies() throws Exception { + OzoneManager om = cluster.getOzoneManager(); + ScmClient scmClient = om.getScmClient(); + String nodeId = + cluster.getStorageContainerManagers().get(0).getSCMNodeId(); + String scmAddrKey = ConfUtils.addKeySuffixes( + OZONE_SCM_ADDRESS_KEY, scmServiceId, nodeId); + + // Update SCM address in OM's live configuration read by proxy providers. + OzoneConfiguration conf = om.getConfiguration(); + conf.set(scmAddrKey, "127.0.0.2"); + + Map changed = new HashMap<>(); + changed.put(scmAddrKey, true); + om.reloadScmProxiesOnReconfig(changed, conf); + + // Both providers must have rebuilt that node's proxy info to the new host. + assertResolvedHost(scmClient.getContainerProxyProvider(), nodeId, + "127.0.0.2"); + assertResolvedHost(scmClient.getBlockProxyProvider(), nodeId, "127.0.0.2"); + } + + /** + * Verifies that proxies are reloaded only when reconfiguration touches SCM node lists or addresses. + */ + @Test + void testReloadScmProxiesOnReconfigIgnoresUnrelatedKey() { + OzoneManager om = cluster.getOzoneManager(); + ScmClient scmClient = om.getScmClient(); + String nodeId = + cluster.getStorageContainerManagers().get(0).getSCMNodeId(); + String scmAddrKey = ConfUtils.addKeySuffixes( + OZONE_SCM_ADDRESS_KEY, scmServiceId, nodeId); + + OzoneConfiguration conf = om.getConfiguration(); + String original = + resolvedHost(scmClient.getContainerProxyProvider(), nodeId); + + // Drift the live address, but report only an unrelated key as changed. + conf.set(scmAddrKey, "127.0.0.9"); + Map changed = new HashMap<>(); + changed.put("ozone.om.unrelated.key", true); + om.reloadScmProxiesOnReconfig(changed, conf); + + // No reload fired, so both providers still resolve the original address. + assertEquals(original, + resolvedHost(scmClient.getContainerProxyProvider(), nodeId)); + assertEquals(original, + resolvedHost(scmClient.getBlockProxyProvider(), nodeId)); + } + + /** + * Verifies that malformed SCM addresses during completion callbacks are caught + * and logged without breaking existing proxies. + */ + @Test + void testReloadScmProxiesOnReconfigCatchesMalformedAddress() { + OzoneManager om = cluster.getOzoneManager(); + ScmClient scmClient = om.getScmClient(); + String nodeId = + cluster.getStorageContainerManagers().get(0).getSCMNodeId(); + String scmAddrKey = ConfUtils.addKeySuffixes( + OZONE_SCM_ADDRESS_KEY, scmServiceId, nodeId); + + OzoneConfiguration conf = om.getConfiguration(); + Set before = + new HashSet<>(scmClient.getContainerProxyProvider().getSCMNodeIds()); + + // Addresses with duplicate ports produce invalid host:port authorities and throw IllegalArgumentException. + conf.set(scmAddrKey, "127.0.0.1:9999"); + Map changed = new HashMap<>(); + changed.put(scmAddrKey, true); + + // The callback must swallow the failure rather than break the chain. + assertDoesNotThrow(() -> om.reloadScmProxiesOnReconfig(changed, conf)); + + // The membership is left in place for both providers. + assertEquals(before, + new HashSet<>(scmClient.getContainerProxyProvider().getSCMNodeIds())); + assertEquals(before, + new HashSet<>(scmClient.getBlockProxyProvider().getSCMNodeIds())); + } + + /** + * Verifies that adding an SCM without a resolved address fails reconfiguration and keeps proxies unchanged. + */ + @Test + void testReconfigureScmNodesFailsWhenAddressMissing() { + OzoneManager om = cluster.getOzoneManager(); + ReconfigurationHandler handler = om.getReconfigurationHandler(); + ScmClient scmClient = om.getScmClient(); + OzoneConfiguration conf = om.getConfiguration(); + String scmNodesKey = + ConfUtils.addKeySuffixes(OZONE_SCM_NODES_KEY, scmServiceId); + + List before = + new ArrayList<>(scmClient.getContainerProxyProvider().getSCMNodeIds()); + String originalValue = conf.get(scmNodesKey); + + List withMissing = new ArrayList<>(before); + withMissing.add("scm-no-address"); + + assertThrows(ReconfigurationException.class, () -> + handler.reconfigureProperty(scmNodesKey, String.join(",", withMissing))); + + // The node list is rolled back and both providers keep their membership. + assertEquals(originalValue, conf.get(scmNodesKey)); + assertEquals(new HashSet<>(before), + new HashSet<>(scmClient.getContainerProxyProvider().getSCMNodeIds())); + assertEquals(new HashSet<>(before), + new HashSet<>(scmClient.getBlockProxyProvider().getSCMNodeIds())); + } + + /** + * Verifies that adding an SCM with a malformed address rolls back node changes and leaves proxies unchanged. + */ + @Test + void testReconfigureScmNodesFailsWhenAddressMalformed() { + OzoneManager om = cluster.getOzoneManager(); + ReconfigurationHandler handler = om.getReconfigurationHandler(); + ScmClient scmClient = om.getScmClient(); + OzoneConfiguration conf = om.getConfiguration(); + String scmNodesKey = + ConfUtils.addKeySuffixes(OZONE_SCM_NODES_KEY, scmServiceId); + String newNodeId = "scm-bad-address"; + String newAddrKey = + ConfUtils.addKeySuffixes(OZONE_SCM_ADDRESS_KEY, scmServiceId, newNodeId); + + List before = + new ArrayList<>(scmClient.getContainerProxyProvider().getSCMNodeIds()); + String originalValue = conf.get(scmNodesKey); + + // Addresses with existing ports form invalid authorities, + // throwing IllegalArgumentException over ConfigurationException. + conf.set(newAddrKey, "127.0.0.1:9999"); + List withBadAddress = new ArrayList<>(before); + withBadAddress.add(newNodeId); + + assertThrows(ReconfigurationException.class, () -> handler + .reconfigureProperty(scmNodesKey, String.join(",", withBadAddress))); + + // The node list is rolled back and both providers keep their membership. + assertEquals(originalValue, conf.get(scmNodesKey)); + assertEquals(new HashSet<>(before), + new HashSet<>(scmClient.getContainerProxyProvider().getSCMNodeIds())); + assertEquals(new HashSet<>(before), + new HashSet<>(scmClient.getBlockProxyProvider().getSCMNodeIds())); + } + + /** + * Verifies that adding an SCM before setting its address rolls back initially, + * then succeeds when retried after setting the address. + */ + @Test + void testReconfigureAddScmNodeNodesBeforeAddress() throws Exception { + OzoneManager om = cluster.getOzoneManager(); + ReconfigurationHandler handler = om.getReconfigurationHandler(); + ScmClient scmClient = om.getScmClient(); + String scmNodesKey = + ConfUtils.addKeySuffixes(OZONE_SCM_NODES_KEY, scmServiceId); + + List before = + new ArrayList<>(scmClient.getContainerProxyProvider().getSCMNodeIds()); + String newNodeId = "scm-added"; + String newAddrKey = + ConfUtils.addKeySuffixes(OZONE_SCM_ADDRESS_KEY, scmServiceId, newNodeId); + List after = new ArrayList<>(before); + after.add(newNodeId); + String requested = String.join(",", after); + + // 1. Unresolvable node list fails and rolls back when applied before the new address. + assertThrows(ReconfigurationException.class, + () -> handler.reconfigureProperty(scmNodesKey, requested)); + assertFalse(new HashSet<>( + scmClient.getContainerProxyProvider().getSCMNodeIds()).contains(newNodeId)); + + // 2. The new SCM's address is applied (prefix key, no reload of its own). + handler.reconfigureProperty(newAddrKey, "127.0.0.1"); + + // 3. Reapplying the node list now resolves the new SCM, so it is added. + handler.reconfigureProperty(scmNodesKey, requested); + Set expected = new HashSet<>(after); + assertEquals(expected, + new HashSet<>(scmClient.getContainerProxyProvider().getSCMNodeIds())); + assertEquals(expected, + new HashSet<>(scmClient.getBlockProxyProvider().getSCMNodeIds())); + } + + private static void assertResolvedHost( + SCMFailoverProxyProviderBase provider, + String nodeId, String expectedHost) { + String host = resolvedHost(provider, nodeId); + assertNotNull(host); + assertEquals(expectedHost, host); + } + + private static String resolvedHost( + SCMFailoverProxyProviderBase provider, String nodeId) { + for (SCMProxyInfo candidate : provider.getSCMProxyInfoList()) { + if (candidate.getNodeId().equals(nodeId)) { + return candidate.getAddress().getAddress().getHostAddress(); + } + } + return null; + } +} diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java index fdc6f95b48db..7901a2a915f5 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java @@ -24,6 +24,8 @@ import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_BLOCK_TOKEN_ENABLED; import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_BLOCK_TOKEN_ENABLED_DEFAULT; import static org.apache.hadoop.hdds.HddsUtils.getScmAddressForClients; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_ADDRESS_KEY; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_NODES_KEY; import static org.apache.hadoop.hdds.server.ServerUtils.updateRPCListenAddress; import static org.apache.hadoop.hdds.utils.HAUtils.getScmInfo; import static org.apache.hadoop.hdds.utils.HddsServerUtil.getRemoteUser; @@ -203,6 +205,8 @@ import org.apache.hadoop.hdds.scm.net.NetworkTopology; import org.apache.hadoop.hdds.scm.protocol.ScmBlockLocationProtocol; import org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocol; +import org.apache.hadoop.hdds.scm.proxy.SCMBlockLocationFailoverProxyProvider; +import org.apache.hadoop.hdds.scm.proxy.SCMContainerLocationFailoverProxyProvider; import org.apache.hadoop.hdds.security.SecurityConfig; import org.apache.hadoop.hdds.security.exception.OzoneSecurityException; import org.apache.hadoop.hdds.security.symmetric.DefaultSecretKeyClient; @@ -252,6 +256,7 @@ import org.apache.hadoop.ozone.audit.OMAction; import org.apache.hadoop.ozone.audit.OMSystemAction; import org.apache.hadoop.ozone.common.Storage.StorageState; +import org.apache.hadoop.ozone.ha.ConfUtils; import org.apache.hadoop.ozone.om.exceptions.OMException; import org.apache.hadoop.ozone.om.exceptions.OMException.ResultCodes; import org.apache.hadoop.ozone.om.exceptions.OMLeaderNotReadyException; @@ -585,6 +590,16 @@ private OzoneManager(OzoneConfiguration conf, StartupOption startupOption) .register(OZONE_READ_BLACKLIST_USERS, this::reconfOzoneReadBlacklistUsers) .register(OZONE_READ_BLACKLIST_GROUPS, this::reconfOzoneReadBlacklistGroups); + // Allow dynamic SCM node/address reconfiguration without restarting OM. + // Address keys use prefix matching to discover new nodes added after startup. + String scmServiceId = HddsUtils.getScmServiceId(conf); + if (scmServiceId != null) { + reconfigurationHandler + .registerPrefix(ConfUtils.addKeySuffixes(OZONE_SCM_ADDRESS_KEY, scmServiceId)) + .register(OZONE_SCM_NODES_KEY + "." + scmServiceId, this::reconfScmNodes) + .registerCompleteCallback(this::reloadScmProxiesOnReconfig); + } + reconfigurationHandler.setReconfigurationCompleteCallback(reconfigurationHandler.defaultLoggingCallback()); reconfigurationHandler.registerCompleteCallback(tracingReconfigurationCallback); @@ -663,12 +678,19 @@ private OzoneManager(OzoneConfiguration conf, StartupOption startupOption) // Honor property 'hadoop.security.token.service.use_ip' omRpcAddressTxt = new Text(SecurityUtil.buildTokenService(omNodeRpcAddr)); - final StorageContainerLocationProtocol scmContainerClient = getScmContainerClient(configuration); + // Retain SCM proxy providers so OM can dynamically reload nodes/addresses without restarting. + final SCMContainerLocationFailoverProxyProvider scmContainerProxyProvider = + new SCMContainerLocationFailoverProxyProvider(configuration, null); + final StorageContainerLocationProtocol scmContainerClient = + HAUtils.getScmContainerClient(configuration, scmContainerProxyProvider); // verifies that the SCM info in the OM Version file is correct. - final ScmBlockLocationProtocol scmBlockClient = getScmBlockClient(configuration); + final SCMBlockLocationFailoverProxyProvider scmBlockProxyProvider = + new SCMBlockLocationFailoverProxyProvider(configuration); + final ScmBlockLocationProtocol scmBlockClient = + HAUtils.getScmBlockClient(configuration, scmBlockProxyProvider); scmTopologyClient = new ScmTopologyClient(scmBlockClient); this.scmClient = new ScmClient(scmBlockClient, scmContainerClient, - configuration); + scmBlockProxyProvider, scmContainerProxyProvider, configuration); this.ozoneLockProvider = new OzoneLockProvider(getKeyPathLockEnabled(), getEnableFileSystemPaths()); @@ -1514,26 +1536,6 @@ private static void loginOMUser(OzoneConfiguration conf) LOG.info("Ozone Manager login successful."); } - /** - * Create a scm block client, used by putKey() and getKey(). - * - * @return {@link ScmBlockLocationProtocol} - */ - private static ScmBlockLocationProtocol getScmBlockClient( - OzoneConfiguration conf) { - return HAUtils.getScmBlockClient(conf); - } - - /** - * Returns a scm container client. - * - * @return {@link StorageContainerLocationProtocol} - */ - private static StorageContainerLocationProtocol getScmContainerClient( - OzoneConfiguration conf) { - return HAUtils.getScmContainerClient(conf); - } - /** * Creates a new instance of rpc server. If an earlier instance is already * running then returns the same. @@ -6092,6 +6094,64 @@ public ListSnapshotDiffJobResponse listSnapshotDiffJobs( } } + /** + * Validates/reloads SCM node list and proxies (block/container) without an OM restart. + * Rolls back on missing node addresses. Requires address configs to be set prior to node list. + */ + private String reconfScmNodes(String value) { + if (StringUtils.isBlank(value)) { + throw new IllegalArgumentException("Reconfiguration failed since setting an empty SCM nodes " + + "configuration is not allowed"); + } + // Publish the node list early so reloadScmNodes() sees it in the live config before the callback finishes. + String scmNodesKey = ConfUtils.addKeySuffixes(OZONE_SCM_NODES_KEY, + HddsUtils.getScmServiceId(configuration)); + String previous = configuration.get(scmNodesKey); + configuration.set(scmNodesKey, value); + try { + scmClient.reloadScmNodes(); + LOG.info("Reloaded SCM proxy configuration for {} : {}", scmNodesKey, value); + } catch (RuntimeException e) { + // Restore previous node list on resolution failure and rethrow to signal FAILED. + if (previous == null) { + configuration.unset(scmNodesKey); + } else { + configuration.set(scmNodesKey, previous); + } + throw e; + } + return value; + } + + /** + * Callback that reloads SCM block and container proxies after reconfiguration finishes. + * Ensures address-only changes take effect and applies newly added nodes if pre-configured. + */ + @VisibleForTesting + public void reloadScmProxiesOnReconfig(Map changedProperties, + Configuration newConf) { + String scmServiceId = HddsUtils.getScmServiceId(configuration); + if (scmServiceId == null || scmClient == null) { + return; + } + String scmNodesKey = ConfUtils.addKeySuffixes(OZONE_SCM_NODES_KEY, scmServiceId); + String scmAddressPrefix = + ConfUtils.addKeySuffixes(OZONE_SCM_ADDRESS_KEY, scmServiceId) + "."; + boolean scmProxyKeyChanged = changedProperties.keySet().stream() + .anyMatch(key -> key.equals(scmNodesKey) || key.startsWith(scmAddressPrefix)); + if (scmProxyKeyChanged) { + try { + scmClient.reloadScmNodes(); + LOG.info("Reloaded SCM failover proxies after reconfiguration of {} / {}*", + scmNodesKey, scmAddressPrefix); + } catch (RuntimeException e) { + // Catch and log reload failures to keep existing proxies and avoid breaking the callback chain. + LOG.warn("Failed to reload SCM failover proxies after reconfiguration of {} / {}*; " + + "keeping the previous SCM proxy configuration", scmNodesKey, scmAddressPrefix, e); + } + } + } + private String reconfOzoneAdmins(String newVal) { Collection admins = OzoneAdmins.getOzoneAdminsFromConfigValue(newVal, omStarterUser); diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ScmClient.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ScmClient.java index a61595faf069..4ec2147111e2 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ScmClient.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/ScmClient.java @@ -24,6 +24,7 @@ import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_CONTAINER_LOCATION_DATANODE_CACHE_SIZE; import static org.apache.hadoop.ozone.om.OMConfigKeys.OZONE_OM_CONTAINER_LOCATION_DATANODE_CACHE_SIZE_DEFAULT; +import com.google.common.annotations.VisibleForTesting; import com.google.common.cache.Cache; import com.google.common.cache.CacheBuilder; import com.google.common.cache.CacheLoader; @@ -46,6 +47,7 @@ import org.apache.hadoop.hdds.scm.pipeline.Pipeline; import org.apache.hadoop.hdds.scm.protocol.ScmBlockLocationProtocol; import org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocol; +import org.apache.hadoop.hdds.scm.proxy.SCMFailoverProxyProviderBase; import org.apache.hadoop.ozone.util.CacheMetrics; /** @@ -55,6 +57,8 @@ public class ScmClient { private final ScmBlockLocationProtocol blockClient; private final StorageContainerLocationProtocol containerClient; + private final SCMFailoverProxyProviderBase blockProxyProvider; + private final SCMFailoverProxyProviderBase containerProxyProvider; private final LoadingCache containerLocationCache; private final CacheMetrics containerCacheMetrics; private final CacheMetrics datanodeDetailsCacheMetrics; @@ -62,8 +66,18 @@ public class ScmClient { ScmClient(ScmBlockLocationProtocol blockClient, StorageContainerLocationProtocol containerClient, OzoneConfiguration configuration) { + this(blockClient, containerClient, null, null, configuration); + } + + ScmClient(ScmBlockLocationProtocol blockClient, + StorageContainerLocationProtocol containerClient, + SCMFailoverProxyProviderBase blockProxyProvider, + SCMFailoverProxyProviderBase containerProxyProvider, + OzoneConfiguration configuration) { this.containerClient = containerClient; this.blockClient = blockClient; + this.blockProxyProvider = blockProxyProvider; + this.containerProxyProvider = containerProxyProvider; Cache datanodeDetailsCache = createDatanodeDetailsCache(configuration); this.containerLocationCache = @@ -144,6 +158,29 @@ static Pipeline newPipelineWithDNCache(Pipeline pipeline, return builder.build(); } + /** + * Reloads block/container SCM proxies on reconfiguration without an OM restart. + * No-op if providers are unavailable (e.g., test mocks). + */ + public void reloadScmNodes() { + if (blockProxyProvider != null) { + blockProxyProvider.changeConfig(); + } + if (containerProxyProvider != null) { + containerProxyProvider.changeConfig(); + } + } + + @VisibleForTesting + public SCMFailoverProxyProviderBase getBlockProxyProvider() { + return blockProxyProvider; + } + + @VisibleForTesting + public SCMFailoverProxyProviderBase getContainerProxyProvider() { + return containerProxyProvider; + } + public ScmBlockLocationProtocol getBlockClient() { return this.blockClient; } diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestScmClient.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestScmClient.java index 89d8c2162c38..b9a08285162d 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestScmClient.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/TestScmClient.java @@ -20,6 +20,7 @@ import static com.google.common.collect.Sets.newHashSet; import static java.util.Arrays.asList; import static org.apache.hadoop.hdds.client.ReplicationConfig.fromTypeAndFactor; +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertThrows; @@ -51,6 +52,7 @@ import org.apache.hadoop.hdds.scm.pipeline.PipelineID; import org.apache.hadoop.hdds.scm.protocol.ScmBlockLocationProtocol; import org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocol; +import org.apache.hadoop.hdds.scm.proxy.SCMFailoverProxyProviderBase; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; @@ -210,6 +212,28 @@ public void testDatanodeDetailsCacheUpdatesIpAddressChange() { assertSame(updated, datanodeDetailsCache.getIfPresent(original.getID())); } + @Test + public void testReloadScmNodesDelegatesToBothProviders() { + SCMFailoverProxyProviderBase blockProvider = + mock(SCMFailoverProxyProviderBase.class); + SCMFailoverProxyProviderBase containerProvider = + mock(SCMFailoverProxyProviderBase.class); + ScmClient client = new ScmClient(mock(ScmBlockLocationProtocol.class), + mock(StorageContainerLocationProtocol.class), blockProvider, + containerProvider, new OzoneConfiguration()); + + client.reloadScmNodes(); + + verify(blockProvider, times(1)).changeConfig(); + verify(containerProvider, times(1)).changeConfig(); + } + + @Test + public void testReloadScmNodesIsNoOpWithoutProviders() { + // Ignore reload when providers are null (e.g., in unit tests). + assertDoesNotThrow(scmClient::reloadScmNodes); + } + ContainerWithPipeline createPipeline(long containerId, List dnList) { ContainerInfo containerInfo = new ContainerInfo.Builder()