This is an automated email from the ASF dual-hosted git repository.

peterxcli pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git


The following commit(s) were added to refs/heads/master by this push:
     new 0f6295ef3ab HDDS-15085. Add DN and Cluster Readiness (#10757)
0f6295ef3ab is described below

commit 0f6295ef3abc613f6bd3d5093a5cb10cde11f3b3
Author: Chun-Hung Tseng <[email protected]>
AuthorDate: Mon Jul 20 19:02:29 2026 +0200

    HDDS-15085. Add DN and Cluster Readiness (#10757)
    
    Co-authored-by: Bolin Lin <[email protected]>
    Co-authored-by: Peter Lee <[email protected]>
---
 .../ozone/local/TestLocalOzoneClusterRuntime.java  |  32 ++-
 hadoop-ozone/tools/pom.xml                         |   4 +
 .../hadoop/ozone/local/LocalOzoneCluster.java      | 228 +++++++++++++++++++--
 .../hadoop/ozone/local/TestLocalOzoneCluster.java  | 109 ++++++++++
 4 files changed, 350 insertions(+), 23 deletions(-)

diff --git 
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneClusterRuntime.java
 
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneClusterRuntime.java
index 2d2d762896a..a4a10f0491b 100644
--- 
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneClusterRuntime.java
+++ 
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneClusterRuntime.java
@@ -17,6 +17,7 @@
 
 package org.apache.hadoop.ozone.local;
 
+import static java.nio.charset.StandardCharsets.UTF_8;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -37,25 +38,28 @@
 import org.junit.jupiter.api.io.TempDir;
 
 /**
- * Integration tests for the SCM and OM portion of {@link LocalOzoneCluster}.
+ * Integration tests for {@link LocalOzoneCluster}.
  */
 class TestLocalOzoneClusterRuntime {
 
+  private static final String KEY_CONTENT = "local ozone key content";
+
   @TempDir
   private Path tempDir;
 
   @Test
-  void scmAndOmStartAndReuseExistingMetadata() throws Exception {
+  void clusterStartsAndReusesExistingData() throws Exception {
     String volumeName = uniqueName("vol");
     String bucketName = uniqueName("bucket");
+    String keyName = uniqueName("key");
     Path dataDir = tempDir.resolve("local-ozone-runtime");
     LocalOzoneClusterConfig config = LocalOzoneClusterConfig.builder(dataDir)
         .setS3gEnabled(false)
         .setStartupTimeout(Duration.ofMinutes(2))
         .build();
 
-    startRuntimeAndCreateBucket(config, volumeName, bucketName);
-    restartRuntimeAndVerifyBucket(config, volumeName, bucketName);
+    startRuntimeAndCreateKey(config, volumeName, bucketName, keyName);
+    restartRuntimeAndVerifyKey(config, volumeName, bucketName, keyName);
   }
 
   @Test
@@ -83,34 +87,44 @@ void formatNeverRejectsUninitializedScmOmStorage() throws 
Exception {
         error.getMessage());
   }
 
-  private void startRuntimeAndCreateBucket(LocalOzoneClusterConfig config,
-      String volumeName, String bucketName) throws Exception {
+  private void startRuntimeAndCreateKey(LocalOzoneClusterConfig config,
+      String volumeName, String bucketName, String keyName) throws Exception {
     try (LocalOzoneCluster cluster = new LocalOzoneCluster(config, new 
OzoneConfiguration())) {
       OzoneConfiguration clientConf =
           cluster.prepareConfiguration().getConfiguration();
       cluster.start();
 
+      assertEquals(config.getDatanodes(), cluster.getDatanodeCount());
       assertServicePortsReachable(cluster);
 
       try (OzoneClient client = OzoneClientFactory.getRpcClient(clientConf)) {
-        TestDataUtil.createVolumeAndBucket(client, volumeName, bucketName);
+        OzoneBucket bucket =
+            TestDataUtil.createVolumeAndBucket(client, volumeName, bucketName);
+        // Writing and reading back a key proves the datanodes registered and
+        // SCM left safe mode, so the cluster is actually usable.
+        TestDataUtil.createKey(bucket, keyName, KEY_CONTENT.getBytes(UTF_8));
+        assertEquals(KEY_CONTENT, TestDataUtil.getKey(bucket, keyName));
       }
     }
   }
 
-  private void restartRuntimeAndVerifyBucket(LocalOzoneClusterConfig config,
-      String volumeName, String bucketName) throws Exception {
+  private void restartRuntimeAndVerifyKey(LocalOzoneClusterConfig config,
+      String volumeName, String bucketName, String keyName) throws Exception {
     try (LocalOzoneCluster cluster = new LocalOzoneCluster(config, new 
OzoneConfiguration())) {
       OzoneConfiguration clientConf =
           cluster.prepareConfiguration().getConfiguration();
       cluster.start();
 
+      assertEquals(config.getDatanodes(), cluster.getDatanodeCount());
       assertServicePortsReachable(cluster);
 
       try (OzoneClient client = OzoneClientFactory.getRpcClient(clientConf)) {
         OzoneVolume volume = client.getObjectStore().getVolume(volumeName);
         OzoneBucket bucket = volume.getBucket(bucketName);
         assertEquals(bucketName, bucket.getName());
+        // Key data written before the restart is still readable from the
+        // persistent datanode storage.
+        assertEquals(KEY_CONTENT, TestDataUtil.getKey(bucket, keyName));
       }
     }
   }
diff --git a/hadoop-ozone/tools/pom.xml b/hadoop-ozone/tools/pom.xml
index bda0ebb3134..3ffab37b2b2 100644
--- a/hadoop-ozone/tools/pom.xml
+++ b/hadoop-ozone/tools/pom.xml
@@ -62,6 +62,10 @@
       <groupId>org.apache.ozone</groupId>
       <artifactId>hdds-config</artifactId>
     </dependency>
+    <dependency>
+      <groupId>org.apache.ozone</groupId>
+      <artifactId>hdds-container-service</artifactId>
+    </dependency>
     <dependency>
       <groupId>org.apache.ozone</groupId>
       <artifactId>hdds-server-framework</artifactId>
diff --git 
a/hadoop-ozone/tools/src/main/java/org/apache/hadoop/ozone/local/LocalOzoneCluster.java
 
b/hadoop-ozone/tools/src/main/java/org/apache/hadoop/ozone/local/LocalOzoneCluster.java
index def3b5b634f..89961e8243a 100644
--- 
a/hadoop-ozone/tools/src/main/java/org/apache/hadoop/ozone/local/LocalOzoneCluster.java
+++ 
b/hadoop-ozone/tools/src/main/java/org/apache/hadoop/ozone/local/LocalOzoneCluster.java
@@ -17,10 +17,16 @@
 
 package org.apache.hadoop.ozone.local;
 
+import static 
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_DATANODE_CLIENT_ADDRESS_KEY;
+import static 
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_DATANODE_CLIENT_BIND_HOST_KEY;
+import static 
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_DATANODE_HTTP_ADDRESS_KEY;
+import static 
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_DATANODE_HTTP_BIND_HOST_KEY;
+import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_HEARTBEAT_INTERVAL;
 import static 
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_SCM_SAFEMODE_MIN_DATANODE;
 import static 
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_SCM_SAFEMODE_PIPELINE_CREATION;
 import static 
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_SCM_WAIT_TIME_AFTER_SAFE_MODE_EXIT;
 import static 
org.apache.hadoop.hdds.scm.ScmConfigKeys.HDDS_CONTAINER_RATIS_ENABLED_KEY;
+import static org.apache.hadoop.hdds.scm.ScmConfigKeys.HDDS_DATANODE_DIR_KEY;
 import static 
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_BLOCK_CLIENT_ADDRESS_KEY;
 import static 
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_BLOCK_CLIENT_BIND_HOST_KEY;
 import static 
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_CLIENT_ADDRESS_KEY;
@@ -41,6 +47,12 @@
 import static 
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_SECURITY_SERVICE_ADDRESS_KEY;
 import static 
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_SECURITY_SERVICE_BIND_HOST_KEY;
 import static org.apache.hadoop.hdds.server.http.BaseHttpServer.SERVER_DIR;
+import static org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_IPC_PORT;
+import static 
org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_RATIS_ADMIN_PORT;
+import static 
org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_RATIS_DATANODE_STORAGE_DIR;
+import static 
org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_RATIS_DATASTREAM_PORT;
+import static 
org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_RATIS_IPC_PORT;
+import static 
org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_RATIS_SERVER_PORT;
 import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_HTTP_BASEDIR;
 import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_METADATA_DIRS;
 import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_REPLICATION;
@@ -68,8 +80,12 @@
 import java.nio.file.Files;
 import java.nio.file.Path;
 import java.time.Duration;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
 import java.util.Comparator;
 import java.util.HashSet;
+import java.util.List;
 import java.util.Objects;
 import java.util.Properties;
 import java.util.Set;
@@ -85,18 +101,19 @@
 import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
 import org.apache.hadoop.hdds.utils.IOUtils;
 import org.apache.hadoop.metrics2.lib.DefaultMetricsSystem;
+import org.apache.hadoop.ozone.HddsDatanodeService;
 import org.apache.hadoop.ozone.OzoneSecurityUtil;
 import org.apache.hadoop.ozone.common.Storage;
+import org.apache.hadoop.ozone.container.replication.ReplicationServer;
 import org.apache.hadoop.ozone.om.OMStorage;
 import org.apache.hadoop.ozone.om.OzoneManager;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 /**
- * Starts the SCM and OM portion of the {@code ozone local} runtime.
+ * Starts the SCM, OM, and datanode portion of the {@code ozone local} runtime.
  *
- * <p>Datanodes, S3 Gateway, Recon, and end-to-end key writes are added by
- * later local runtime tickets.</p>
+ * <p>S3 Gateway and Recon are added by later local runtime tickets.</p>
  */
 public final class LocalOzoneCluster implements LocalOzoneRuntime {
 
@@ -111,6 +128,7 @@ public final class LocalOzoneCluster implements 
LocalOzoneRuntime {
   private static final String OZONE_METADATA_DIR_NAME = "ozone-metadata";
   private static final String DATA_DIR_NAME = "data";
   private static final String RATIS_DIR_NAME = "ratis";
+  private static final String DATANODE_DIR_PREFIX = "datanode-";
   private static final String SCM_CLIENT_PORT_KEY = "scm.client";
   private static final String SCM_BLOCK_PORT_KEY = "scm.block";
   private static final String SCM_DATANODE_PORT_KEY = "scm.datanode";
@@ -123,6 +141,15 @@ public final class LocalOzoneCluster implements 
LocalOzoneRuntime {
   private static final String OM_HTTP_PORT_KEY = "om.http";
   private static final String OM_HTTPS_PORT_KEY = "om.https";
   private static final String OM_RATIS_PORT_KEY = "om.ratis";
+  private static final String DATANODE_PORT_KEY_PREFIX = "dn.";
+  private static final String DATANODE_HTTP_PORT_KEY_SUFFIX = "http";
+  private static final String DATANODE_CLIENT_PORT_KEY_SUFFIX = "client";
+  private static final String DATANODE_CONTAINER_IPC_PORT_KEY_SUFFIX = 
"container.ipc";
+  private static final String DATANODE_RATIS_IPC_PORT_KEY_SUFFIX = "ratis.ipc";
+  private static final String DATANODE_RATIS_ADMIN_PORT_KEY_SUFFIX = 
"ratis.admin";
+  private static final String DATANODE_RATIS_SERVER_PORT_KEY_SUFFIX = 
"ratis.server";
+  private static final String DATANODE_RATIS_DATASTREAM_PORT_KEY_SUFFIX = 
"ratis.datastream";
+  private static final String DATANODE_REPLICATION_PORT_KEY_SUFFIX = 
"replication";
   private static final int LOCAL_RATIS_RPC_TIMEOUT_SECONDS = 1;
   private static final long SCM_CLIENT_MAX_RETRY_TIMEOUT_MILLIS = 30_000;
   private static final long READINESS_POLL_INTERVAL_MILLIS = 500;
@@ -144,6 +171,23 @@ public final class LocalOzoneCluster implements 
LocalOzoneRuntime {
       OM_RATIS_PORT_KEY
   };
 
+  private static final String[] DATANODE_PORT_KEY_SUFFIXES = {
+      DATANODE_HTTP_PORT_KEY_SUFFIX,
+      DATANODE_CLIENT_PORT_KEY_SUFFIX,
+      DATANODE_CONTAINER_IPC_PORT_KEY_SUFFIX,
+      DATANODE_RATIS_IPC_PORT_KEY_SUFFIX,
+      DATANODE_RATIS_ADMIN_PORT_KEY_SUFFIX,
+      DATANODE_RATIS_SERVER_PORT_KEY_SUFFIX,
+      DATANODE_RATIS_DATASTREAM_PORT_KEY_SUFFIX,
+      DATANODE_REPLICATION_PORT_KEY_SUFFIX
+  };
+
+  // Every datanode runs in this JVM and reserves 
DATANODE_PORT_KEY_SUFFIXES.length
+  // local ports, so an unbounded count would exhaust local ports; cap it.
+  static final int MAX_DATANODES = 20;
+
+  private static final String[] NO_ARGS = new String[0];
+
   private final LocalOzoneClusterConfig config;
   private final OzoneConfiguration seedConfiguration;
   private boolean closed;
@@ -151,6 +195,7 @@ public final class LocalOzoneCluster implements 
LocalOzoneRuntime {
   private PreparedConfiguration preparedConfiguration;
   private StorageContainerManager scm;
   private OzoneManager om;
+  private final List<HddsDatanodeService> datanodes = new ArrayList<>();
   private boolean previousMetricsMiniClusterMode;
   private boolean metricsMiniClusterModeEnabled;
 
@@ -179,7 +224,8 @@ public void start() throws Exception {
       initializeStorage(prepared.getConfiguration());
       startScm(prepared.getConfiguration());
       startOm(prepared.getConfiguration());
-      waitForScmAndOmReadiness(config.getStartupTimeout());
+      startDatanodes(prepared.getDatanodeConfigurations());
+      waitForClusterReadiness(config.getStartupTimeout());
     } catch (Exception ex) {
       // Roll back without latching closed: the caller's close() still owns the
       // ephemeral data dir lifecycle.
@@ -202,9 +248,12 @@ PreparedConfiguration prepareConfiguration() throws 
IOException {
     PortAllocator portAllocator = new PortAllocator();
     int scmPort = configureScm(conf, persistedPorts, portAllocator);
     int omPort = configureOm(conf, persistedPorts, portAllocator);
+    List<OzoneConfiguration> datanodeConfigurations =
+        configureDatanodes(conf, persistedPorts, portAllocator);
 
     persistedPorts.store();
-    preparedConfiguration = new PreparedConfiguration(conf, scmPort, omPort);
+    preparedConfiguration = new PreparedConfiguration(conf, scmPort, omPort,
+        datanodeConfigurations);
     return preparedConfiguration;
   }
 
@@ -224,6 +273,13 @@ public int getOmPort() {
     return om.getOmRpcServerAddr().getPort();
   }
 
+  /**
+   * Returns the number of running datanodes.
+   */
+  public int getDatanodeCount() {
+    return datanodes.size();
+  }
+
   @Override
   public int getS3gPort() {
     return -1;
@@ -258,13 +314,26 @@ public void close() throws IOException {
 
   private void stopServices() {
     try {
-      // Shutdown is best-effort so one failed service cannot leak the other.
-      IOUtils.closeQuietly(this::stopOm, this::stopScm);
+      // Shutdown is best-effort so one failed service cannot leak the others.
+      IOUtils.closeQuietly(this::stopDatanodes, this::stopOm, this::stopScm);
     } finally {
       restoreSameJvmMetricsMode();
     }
   }
 
+  private void stopDatanodes() {
+    List<AutoCloseable> stoppers = new ArrayList<>();
+    for (int i = datanodes.size() - 1; i >= 0; i--) {
+      HddsDatanodeService service = datanodes.get(i);
+      stoppers.add(() -> {
+        service.stop();
+        service.join();
+      });
+    }
+    datanodes.clear();
+    IOUtils.closeQuietly(stoppers);
+  }
+
   private void stopOm() {
     OzoneManager service = om;
     om = null;
@@ -293,6 +362,11 @@ private void configureLocalDefaults(OzoneConfiguration 
conf) {
     conf.set(OZONE_SERVER_DEFAULT_REPLICATION_TYPE_KEY,
         ReplicationType.STAND_ALONE.name());
     conf.setBoolean(HDDS_CONTAINER_RATIS_ENABLED_KEY, false);
+    // A single-node local cluster can heartbeat aggressively; this speeds
+    // datanode registration and safe-mode exit. Use set(), not setIfUnset():
+    // ozone-default.xml supplies the 30s default that would otherwise defeat
+    // the override.
+    conf.set(HDDS_HEARTBEAT_INTERVAL, "1s");
     conf.setBoolean(HDDS_SCM_SAFEMODE_PIPELINE_CREATION, false);
     conf.setInt(HDDS_SCM_SAFEMODE_MIN_DATANODE,
         Math.max(1, config.getDatanodes()));
@@ -403,6 +477,81 @@ private void configureOmStorage(OzoneConfiguration conf) 
throws IOException {
         omMetadataDir.toString());
   }
 
+  private List<OzoneConfiguration> configureDatanodes(OzoneConfiguration conf,
+      PersistedPortState persistedPorts, PortAllocator portAllocator)
+      throws IOException {
+    int datanodeCount = config.getDatanodes();
+    if (datanodeCount > MAX_DATANODES) {
+      throw new IOException("Datanode count " + datanodeCount
+          + " exceeds the local maximum of " + MAX_DATANODES
+          + "; each datanode reserves " + DATANODE_PORT_KEY_SUFFIXES.length
+          + " local ports.");
+    }
+    List<OzoneConfiguration> datanodeConfigurations =
+        new ArrayList<>(config.getDatanodes());
+    for (int index = 0; index < config.getDatanodes(); index++) {
+      datanodeConfigurations.add(
+          configureDatanode(conf, index, persistedPorts, portAllocator));
+    }
+    return datanodeConfigurations;
+  }
+
+  private OzoneConfiguration configureDatanode(OzoneConfiguration conf,
+      int index, PersistedPortState persistedPorts,
+      PortAllocator portAllocator) throws IOException {
+    OzoneConfiguration dnConf = new OzoneConfiguration(conf);
+    configureDatanodeStorage(dnConf, index);
+
+    dnConf.set(HDDS_DATANODE_HTTP_ADDRESS_KEY, address(config.getHost(),
+        reserveDatanodePort(portAllocator, persistedPorts, index,
+            DATANODE_HTTP_PORT_KEY_SUFFIX)));
+    dnConf.set(HDDS_DATANODE_HTTP_BIND_HOST_KEY, config.getBindHost());
+    dnConf.set(HDDS_DATANODE_CLIENT_ADDRESS_KEY, address(config.getHost(),
+        reserveDatanodePort(portAllocator, persistedPorts, index,
+            DATANODE_CLIENT_PORT_KEY_SUFFIX)));
+    dnConf.set(HDDS_DATANODE_CLIENT_BIND_HOST_KEY, config.getBindHost());
+    dnConf.setInt(HDDS_CONTAINER_IPC_PORT,
+        reserveDatanodePort(portAllocator, persistedPorts, index,
+            DATANODE_CONTAINER_IPC_PORT_KEY_SUFFIX));
+    dnConf.setInt(HDDS_CONTAINER_RATIS_IPC_PORT,
+        reserveDatanodePort(portAllocator, persistedPorts, index,
+            DATANODE_RATIS_IPC_PORT_KEY_SUFFIX));
+    dnConf.setInt(HDDS_CONTAINER_RATIS_ADMIN_PORT,
+        reserveDatanodePort(portAllocator, persistedPorts, index,
+            DATANODE_RATIS_ADMIN_PORT_KEY_SUFFIX));
+    dnConf.setInt(HDDS_CONTAINER_RATIS_SERVER_PORT,
+        reserveDatanodePort(portAllocator, persistedPorts, index,
+            DATANODE_RATIS_SERVER_PORT_KEY_SUFFIX));
+    dnConf.setInt(HDDS_CONTAINER_RATIS_DATASTREAM_PORT,
+        reserveDatanodePort(portAllocator, persistedPorts, index,
+            DATANODE_RATIS_DATASTREAM_PORT_KEY_SUFFIX));
+
+    ReplicationServer.ReplicationConfig replicationConfig =
+        dnConf.getObject(ReplicationServer.ReplicationConfig.class);
+    replicationConfig.setPort(reserveDatanodePort(portAllocator,
+        persistedPorts, index, DATANODE_REPLICATION_PORT_KEY_SUFFIX));
+    dnConf.setFromObject(replicationConfig);
+    return dnConf;
+  }
+
+  private void configureDatanodeStorage(OzoneConfiguration dnConf, int index)
+      throws IOException {
+    Path datanodeDir = config.getDataDir()
+        .resolve(DATANODE_DIR_PREFIX + index);
+    Path datanodeMetadataDir = datanodeDir.resolve(OZONE_METADATA_DIR_NAME);
+    Files.createDirectories(datanodeMetadataDir);
+    Files.createDirectories(datanodeDir.resolve(DATA_DIR_NAME));
+
+    dnConf.set(OZONE_METADATA_DIRS, datanodeMetadataDir.toString());
+    // Each datanode gets its own Jetty base dir so same-JVM HTTP servers do
+    // not share unpacked web resources.
+    dnConf.set(OZONE_HTTP_BASEDIR, datanodeMetadataDir + SERVER_DIR);
+    dnConf.set(HDDS_DATANODE_DIR_KEY,
+        datanodeDir.resolve(DATA_DIR_NAME).toString());
+    dnConf.set(HDDS_CONTAINER_RATIS_DATANODE_STORAGE_DIR,
+        datanodeDir.resolve(RATIS_DIR_NAME).toString());
+  }
+
   private void initializeStorage(OzoneConfiguration conf) throws IOException {
     SCMStorageConfig scmStorage = new SCMStorageConfig(conf);
     OMStorage omStorage = new OMStorage(conf);
@@ -483,22 +632,43 @@ private void startOm(OzoneConfiguration conf) throws 
Exception {
     om.start();
   }
 
-  private void waitForScmAndOmReadiness(Duration timeout) throws Exception {
+  private void startDatanodes(List<OzoneConfiguration> datanodeConfigurations) 
{
+    for (OzoneConfiguration dnConf : datanodeConfigurations) {
+      // Track the datanode before start() so a failed start can still be
+      // rolled back by stopServices().
+      HddsDatanodeService datanode = new HddsDatanodeService(NO_ARGS);
+      datanodes.add(datanode);
+      datanode.start(dnConf);
+    }
+  }
+
+  private void waitForClusterReadiness(Duration timeout) throws Exception {
     long deadlineNanos = System.nanoTime() + timeout.toNanos();
     while (true) {
-      // Readiness is intentionally scoped to SCM and OM leadership until
-      // datanode and full safe-mode readiness are added in later tickets.
-      if (scm.checkLeader() && om.isLeaderReady()) {
+      if (isClusterReady()) {
         return;
       }
       if (System.nanoTime() >= deadlineNanos) {
         throw new TimeoutException("Timed out waiting " + timeout
-            + " for local SCM and OM leadership.");
+            + " for the local Ozone cluster to become ready.");
       }
       Thread.sleep(READINESS_POLL_INTERVAL_MILLIS);
     }
   }
 
+  private boolean isClusterReady() {
+    if (!scm.checkLeader() || !om.isLeaderReady()) {
+      return false;
+    }
+    if (config.getDatanodes() == 0) {
+      return true;
+    }
+    // The cluster is usable once every datanode has registered with SCM and
+    // SCM has left safe mode.
+    return scm.getScmNodeManager().getAllNodes().size() >= 
config.getDatanodes()
+        && !scm.isInSafeMode();
+  }
+
   private void enableSameJvmMetricsMode() {
     if (!metricsMiniClusterModeEnabled) {
       previousMetricsMiniClusterMode = 
DefaultMetricsSystem.inMiniClusterMode();
@@ -567,11 +737,22 @@ private PersistedPortState loadPersistedPortState() 
throws IOException {
     PersistedPortState persistedPorts =
         PersistedPortState.load(portStateFile());
     if (config.getFormatMode() == LocalOzoneClusterConfig.FormatMode.NEVER) {
-      persistedPorts.requireKeys(REQUIRED_PERSISTED_PORT_KEYS);
+      persistedPorts.requireKeys(requiredPersistedPortKeys());
     }
     return persistedPorts;
   }
 
+  private String[] requiredPersistedPortKeys() {
+    List<String> keys =
+        new ArrayList<>(Arrays.asList(REQUIRED_PERSISTED_PORT_KEYS));
+    for (int index = 0; index < config.getDatanodes(); index++) {
+      for (String suffix : DATANODE_PORT_KEY_SUFFIXES) {
+        keys.add(datanodePortKey(index, suffix));
+      }
+    }
+    return keys.toArray(new String[0]);
+  }
+
   private int reservePort(PortAllocator allocator,
       PersistedPortState persistedPorts, String key, int configuredPort)
       throws IOException {
@@ -582,6 +763,17 @@ private int reservePort(PortAllocator allocator,
     return port;
   }
 
+  private int reserveDatanodePort(PortAllocator allocator,
+      PersistedPortState persistedPorts, int index, String suffix)
+      throws IOException {
+    return reservePort(allocator, persistedPorts,
+        datanodePortKey(index, suffix), 0);
+  }
+
+  private static String datanodePortKey(int index, String suffix) {
+    return DATANODE_PORT_KEY_PREFIX + index + "." + suffix;
+  }
+
   private Path metadataDir() {
     return config.getDataDir().resolve(METADATA_DIR_NAME);
   }
@@ -615,13 +807,17 @@ static final class PreparedConfiguration {
     private final OzoneConfiguration configuration;
     private final int scmPort;
     private final int omPort;
+    private final List<OzoneConfiguration> datanodeConfigurations;
 
     PreparedConfiguration(OzoneConfiguration configuration, int scmPort,
-        int omPort) {
+        int omPort, List<OzoneConfiguration> datanodeConfigurations) {
       this.configuration = Objects.requireNonNull(configuration,
           "configuration");
       this.scmPort = scmPort;
       this.omPort = omPort;
+      this.datanodeConfigurations = Collections.unmodifiableList(
+          new ArrayList<>(Objects.requireNonNull(datanodeConfigurations,
+              "datanodeConfigurations")));
     }
 
     OzoneConfiguration getConfiguration() {
@@ -635,6 +831,10 @@ int getScmPort() {
     int getOmPort() {
       return omPort;
     }
+
+    List<OzoneConfiguration> getDatanodeConfigurations() {
+      return datanodeConfigurations;
+    }
   }
 
   /**
diff --git 
a/hadoop-ozone/tools/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneCluster.java
 
b/hadoop-ozone/tools/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneCluster.java
index 8ecc3a9db98..1f109c6de42 100644
--- 
a/hadoop-ozone/tools/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneCluster.java
+++ 
b/hadoop-ozone/tools/src/test/java/org/apache/hadoop/ozone/local/TestLocalOzoneCluster.java
@@ -21,8 +21,10 @@
 import static 
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_SCM_SAFEMODE_MIN_DATANODE;
 import static 
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_SCM_SAFEMODE_PIPELINE_CREATION;
 import static 
org.apache.hadoop.hdds.scm.ScmConfigKeys.HDDS_CONTAINER_RATIS_ENABLED_KEY;
+import static org.apache.hadoop.hdds.scm.ScmConfigKeys.HDDS_DATANODE_DIR_KEY;
 import static 
org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_CLIENT_ADDRESS_KEY;
 import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_NAMES;
+import static org.apache.hadoop.ozone.OzoneConfigKeys.HDDS_CONTAINER_IPC_PORT;
 import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_METADATA_DIRS;
 import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_REPLICATION;
 import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_REPLICATION_TYPE;
@@ -305,6 +307,113 @@ void persistedPortFileContainsDistinctAllocatedPorts() 
throws Exception {
         properties.getProperty("om.rpc"));
   }
 
+  @Test
+  void prepareConfigurationCreatesDatanodeConfigurations() throws Exception {
+    Path dataDir = tempDir.resolve("local-ozone");
+    LocalOzoneClusterConfig config = LocalOzoneClusterConfig.builder(dataDir)
+        .setDatanodes(2)
+        .build();
+
+    LocalOzoneCluster.PreparedConfiguration prepared = prepare(config);
+
+    assertEquals(2, prepared.getDatanodeConfigurations().size());
+    for (int index = 0; index < 2; index++) {
+      OzoneConfiguration dnConf = 
prepared.getDatanodeConfigurations().get(index);
+      Path datanodeDir = dataDir.resolve("datanode-" + index);
+      assertTrue(Files.isDirectory(datanodeDir.resolve("ozone-metadata")));
+      assertTrue(Files.isDirectory(datanodeDir.resolve("data")));
+      assertEquals(datanodeDir.resolve("ozone-metadata").toString(),
+          dnConf.get(OZONE_METADATA_DIRS));
+      assertEquals(datanodeDir.resolve("data").toString(),
+          dnConf.get(HDDS_DATANODE_DIR_KEY));
+      assertTrue(dnConf.getInt(HDDS_CONTAINER_IPC_PORT, 0) > 0);
+    }
+    assertNotEquals(
+        prepared.getDatanodeConfigurations().get(0)
+            .getInt(HDDS_CONTAINER_IPC_PORT, 0),
+        prepared.getDatanodeConfigurations().get(1)
+            .getInt(HDDS_CONTAINER_IPC_PORT, 0));
+  }
+
+  @Test
+  void prepareConfigurationRejectsTooManyDatanodes() throws Exception {
+    LocalOzoneClusterConfig config = LocalOzoneClusterConfig.builder(
+            tempDir.resolve("local-ozone"))
+        .setDatanodes(LocalOzoneCluster.MAX_DATANODES + 1)
+        .build();
+
+    IOException error = assertPrepareFails(config);
+
+    assertEquals("Datanode count " + (LocalOzoneCluster.MAX_DATANODES + 1)
+        + " exceeds the local maximum of " + LocalOzoneCluster.MAX_DATANODES
+        + "; each datanode reserves 8 local ports.", error.getMessage());
+  }
+
+  @Test
+  void persistedPortFileContainsDatanodePorts() throws Exception {
+    Path dataDir = tempDir.resolve("local-ozone");
+    LocalOzoneClusterConfig config = LocalOzoneClusterConfig.builder(dataDir)
+        .setDatanodes(1)
+        .build();
+
+    prepare(config);
+
+    Properties properties = loadPortState(dataDir);
+    assertPositivePort(properties, "dn.0.http");
+    assertPositivePort(properties, "dn.0.client");
+    assertPositivePort(properties, "dn.0.container.ipc");
+    assertPositivePort(properties, "dn.0.ratis.ipc");
+    assertPositivePort(properties, "dn.0.ratis.admin");
+    assertPositivePort(properties, "dn.0.ratis.server");
+    assertPositivePort(properties, "dn.0.ratis.datastream");
+    assertPositivePort(properties, "dn.0.replication");
+  }
+
+  @Test
+  void prepareConfigurationPersistsDatanodePortsAcrossInstances()
+      throws Exception {
+    Path dataDir = tempDir.resolve("local-ozone");
+    LocalOzoneClusterConfig config =
+        LocalOzoneClusterConfig.builder(dataDir).build();
+
+    LocalOzoneCluster.PreparedConfiguration first = prepare(config);
+    LocalOzoneCluster.PreparedConfiguration second = prepare(config);
+
+    assertEquals(
+        first.getDatanodeConfigurations().get(0)
+            .getInt(HDDS_CONTAINER_IPC_PORT, 0),
+        second.getDatanodeConfigurations().get(0)
+            .getInt(HDDS_CONTAINER_IPC_PORT, 0));
+  }
+
+  @Test
+  void formatNeverRejectsPortStateMissingDatanodePorts() throws Exception {
+    Path dataDir = tempDir.resolve("local-ozone");
+    prepare(LocalOzoneClusterConfig.builder(dataDir).setDatanodes(1).build());
+    LocalOzoneClusterConfig config = LocalOzoneClusterConfig.builder(dataDir)
+        .setFormatMode(LocalOzoneClusterConfig.FormatMode.NEVER)
+        .setDatanodes(2)
+        .build();
+
+    IOException error = assertPrepareFails(config);
+
+    assertMessageContains(error, "dn.1.");
+  }
+
+  @Test
+  void getDatanodeCountReturnsZeroBeforeStart() throws Exception {
+    LocalOzoneClusterConfig config = LocalOzoneClusterConfig.builder(
+            tempDir.resolve("local-ozone"))
+        .setDatanodes(3)
+        .build();
+
+    try (LocalOzoneCluster cluster = newCluster(config)) {
+      cluster.prepareConfiguration();
+
+      assertEquals(0, cluster.getDatanodeCount());
+    }
+  }
+
   private LocalOzoneCluster.PreparedConfiguration prepare(
       LocalOzoneClusterConfig config) throws IOException {
     try (LocalOzoneCluster cluster = newCluster(config)) {


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to