diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/EndpointStateMachine.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/EndpointStateMachine.java index 60bf26a2dc0f..8e98705cade1 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/EndpointStateMachine.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/statemachine/EndpointStateMachine.java @@ -49,7 +49,9 @@ public class EndpointStateMachine private final HostAndPort hostAndPort; private final Lock lock; private final ConfigurationSource conf; - private EndPointStates state = EndPointStates.FIRST; + // RunningDatanodeState reads this without the endpoint lock and must see late SHUTDOWN transitions, + // even after its wait for the endpoint task has timed out. + private volatile EndPointStates state = EndPointStates.FIRST; private VersionResponse version; private ZonedDateTime lastSuccessfulHeartbeat; private boolean isPassive; diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/datanode/RunningDatanodeState.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/datanode/RunningDatanodeState.java index 6b8ea712a93d..fb7eac48db6d 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/datanode/RunningDatanodeState.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/datanode/RunningDatanodeState.java @@ -52,9 +52,7 @@ public class RunningDatanodeState implements DatanodeState { private final ConfigurationSource conf; private final StateContext context; private CompletionService ecs; - // Since we connectionManager endpoints can be changed by reconfiguration - // we should not rely on ConnectionManager#getValues being unchanged between - // execute and await + // Include tasks from earlier heartbeats whose completions have not been collected yet. private int executingEndpointCount = 0; public RunningDatanodeState(ConfigurationSource conf, @@ -88,8 +86,10 @@ public void onExit() { */ @Override public void execute(ExecutorService executor) { - ecs = new ExecutorCompletionService<>(executor); - executingEndpointCount = 0; + // Reuse the completion queue across heartbeats so results arriving after await() times out are not lost. + if (ecs == null) { + ecs = new ExecutorCompletionService<>(executor); + } for (EndpointStateMachine endpoint : connectionManager.getValues()) { Callable endpointTask = buildEndPointTask(endpoint); if (endpointTask != null) { @@ -210,26 +210,21 @@ private Callable buildEndPointTask( public DatanodeStateMachine.DatanodeStates await(long duration, TimeUnit timeUnit) throws InterruptedException { - int returned = 0; long durationMS = timeUnit.toMillis(duration); long timeLeft = durationMS; long startTime = Time.monotonicNow(); List> results = new LinkedList<>(); - while (returned < executingEndpointCount && timeLeft > 0) { + while (executingEndpointCount > 0 && timeLeft > 0) { Future result = ecs.poll(timeLeft, TimeUnit.MILLISECONDS); if (result != null) { results.add(result); - returned++; + executingEndpointCount--; } timeLeft = durationMS - (Time.monotonicNow() - startTime); } return computeNextContainerState(results); } - @Override - public void clear() { - ecs = null; - } } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/endpoint/VersionEndpointTask.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/endpoint/VersionEndpointTask.java index 40844c563e9a..a1e9cbc3765f 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/endpoint/VersionEndpointTask.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/endpoint/VersionEndpointTask.java @@ -18,7 +18,6 @@ package org.apache.hadoop.ozone.container.common.states.endpoint; import java.io.IOException; -import java.net.BindException; import java.util.Objects; import java.util.concurrent.Callable; import org.apache.hadoop.hdds.conf.ConfigurationSource; @@ -72,20 +71,27 @@ public EndpointStateMachine.EndPointStates call() throws Exception { if (!rpcEndPoint.isPassive()) { // If end point is passive, datanode does not need to check volumes. - String scmId = response.getValue(OzoneConsts.SCM_ID); - String clusterId = response.getValue(OzoneConsts.CLUSTER_ID); + try { + String scmId = response.getValue(OzoneConsts.SCM_ID); + String clusterId = response.getValue(OzoneConsts.CLUSTER_ID); - Objects.requireNonNull(scmId, "scmId == null"); - Objects.requireNonNull(clusterId, "clusterId == null"); + Objects.requireNonNull(scmId, "scmId == null"); + Objects.requireNonNull(clusterId, "clusterId == null"); - // Check DbVolumes, format DbVolume at first register time. - checkVolumeSet(ozoneContainer.getDbVolumeSet(), scmId, clusterId); + // Check DbVolumes, format DbVolume at first register time. + checkVolumeSet(ozoneContainer.getDbVolumeSet(), scmId, clusterId); - // Check HddsVolumes - checkVolumeSet(ozoneContainer.getVolumeSet(), scmId, clusterId); + // Check HddsVolumes + checkVolumeSet(ozoneContainer.getVolumeSet(), scmId, clusterId); - // Start the container services after getting the version information - ozoneContainer.start(clusterId); + // Start the container services after getting the version information + ozoneContainer.start(clusterId); + } catch (Exception | Error ex) { + // Handle this in the task: its caller may already have timed out waiting for startup. + LOG.error("Failed to start required container services for SCM {}. Shutting down datanode.", + rpcEndPoint.getAddress(), ex); + return rpcEndPoint.setState(EndpointStateMachine.EndPointStates.SHUTDOWN); + } } EndpointStateMachine.EndPointStates nextState = rpcEndPoint.getState().getNextState(); @@ -95,9 +101,8 @@ public EndpointStateMachine.EndPointStates call() throws Exception { LOG.debug("Cannot execute GetVersion task as endpoint state machine " + "is in {} state", rpcEndPoint.getState()); } - } catch (DiskOutOfSpaceException | BindException ex) { - rpcEndPoint.setState(EndpointStateMachine.EndPointStates.SHUTDOWN); } catch (IOException ex) { + // Communication failures are retryable; local initialization failures are handled above. rpcEndPoint.logIfNeeded(ex); } finally { rpcEndPoint.unlock(); diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/ozoneimpl/OzoneContainer.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/ozoneimpl/OzoneContainer.java index 08d5c18c1d15..ff1fedaa6ace 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/ozoneimpl/OzoneContainer.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/ozoneimpl/OzoneContainer.java @@ -51,7 +51,6 @@ import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicLong; -import java.util.concurrent.atomic.AtomicReference; import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.conf.ConfigurationSource; import org.apache.hadoop.hdds.protocol.DatanodeDetails; @@ -147,7 +146,10 @@ public class OzoneContainer { recoveringContainerScrubbingService; private final GrpcTlsConfig tlsClientConfig; private DiskBalancerService diskBalancerService; - private final AtomicReference initializingStatus; + private final Object initializationLock = new Object(); + // Guarded by initializationLock. + private InitializingStatus initializingStatus; + private Throwable initializationFailure; private final ReplicationServer replicationServer; private DatanodeDetails datanodeDetails; private StateContext context; @@ -160,7 +162,7 @@ public class OzoneContainer { private final DatanodeStorageMetrics datanodeStorageMetrics; enum InitializingStatus { - UNINITIALIZED, INITIALIZING, INITIALIZED + UNINITIALIZED, INITIALIZING, INITIALIZED, FAILED } /** @@ -334,7 +336,7 @@ public OzoneContainer(HddsDatanodeService hddsDatanodeService, datanodeStorageMetrics = DatanodeStorageMetrics.create(volumeSet); - initializingStatus = new AtomicReference<>(InitializingStatus.UNINITIALIZED); + initializingStatus = InitializingStatus.UNINITIALIZED; } /** @@ -547,24 +549,30 @@ public OnDemandContainerScanner getOnDemandScanner() { * @throws IOException */ public void start(String clusterId) throws IOException { - // If SCM HA is enabled, OzoneContainer#start() will be called multi-times - // from VersionEndpointTask. The first call should do the initializing job, - // the successive calls should wait until OzoneContainer is initialized. - if (!initializingStatus.compareAndSet( - InitializingStatus.UNINITIALIZED, InitializingStatus.INITIALIZING)) { - - // wait OzoneContainer to finish its initializing. - while (initializingStatus.get() != InitializingStatus.INITIALIZED) { - try { - Thread.sleep(1); - } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } + synchronized (initializationLock) { + // SCM endpoints share one initialization attempt, including its failure. + if (initializingStatus == InitializingStatus.INITIALIZED) { + LOG.info("Ignore. OzoneContainer already started."); + return; + } + if (initializingStatus == InitializingStatus.FAILED) { + throw new IOException("OzoneContainer initialization previously failed", initializationFailure); + } + + initializingStatus = InitializingStatus.INITIALIZING; + try { + initializeContainerServices(clusterId); + initializingStatus = InitializingStatus.INITIALIZED; + } catch (IOException | RuntimeException | Error ex) { + // Partially started services cannot safely be initialized again. + initializationFailure = ex; + initializingStatus = InitializingStatus.FAILED; + throw ex; } - LOG.info("Ignore. OzoneContainer already started."); - return; } + } + private void initializeContainerServices(String clusterId) throws IOException { DatanodeLayoutStorage layoutStorage = new DatanodeLayoutStorage(config); layoutStorage.setClusterId(clusterId); @@ -607,9 +615,6 @@ public void start(String clusterId) throws IOException { recoveringContainerScrubbingService.start(); initHddsVolumeContainer(); - - // mark OzoneContainer as INITIALIZED. - initializingStatus.set(InitializingStatus.INITIALIZED); } /** diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/TestDatanodeStateMachine.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/TestDatanodeStateMachine.java index 721f50bec139..ef03aa148dff 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/TestDatanodeStateMachine.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/TestDatanodeStateMachine.java @@ -18,44 +18,65 @@ package org.apache.hadoop.ozone.container.common; import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_HEARTBEAT_RPC_TIMEOUT; +import static org.assertj.core.api.Assertions.assertThat; import static org.junit.jupiter.api.Assertions.assertEquals; 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 static org.junit.jupiter.api.Assertions.fail; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import com.google.common.collect.Maps; import com.google.common.util.concurrent.ThreadFactoryBuilder; import java.io.File; import java.io.IOException; +import java.lang.reflect.Field; +import java.net.Socket; import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.UUID; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.ExecutorService; +import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import org.apache.hadoop.hdds.conf.OzoneConfiguration; +import org.apache.hadoop.hdds.conf.ReconfigurationHandler; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.scm.ScmConfigKeys; import org.apache.hadoop.hdds.upgrade.HDDSLayoutFeature; import org.apache.hadoop.ipc_.RPC; +import org.apache.hadoop.ozone.HddsDatanodeStopService; import org.apache.hadoop.ozone.OzoneConfigKeys; import org.apache.hadoop.ozone.OzoneConsts; import org.apache.hadoop.ozone.container.common.helpers.ContainerUtils; import org.apache.hadoop.ozone.container.common.interfaces.VolumeChoosingPolicy; import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine.DatanodeStates; import org.apache.hadoop.ozone.container.common.statemachine.EndpointStateMachine; import org.apache.hadoop.ozone.container.common.statemachine.SCMConnectionManager; import org.apache.hadoop.ozone.container.common.states.DatanodeState; import org.apache.hadoop.ozone.container.common.states.datanode.InitDatanodeState; import org.apache.hadoop.ozone.container.common.states.datanode.RunningDatanodeState; +import org.apache.hadoop.ozone.container.common.states.endpoint.VersionEndpointTask; +import org.apache.hadoop.ozone.container.common.transport.server.XceiverServerSpi; import org.apache.hadoop.ozone.container.common.volume.CapacityVolumeChoosingPolicy; +import org.apache.hadoop.ozone.container.ozoneimpl.OzoneContainer; +import org.apache.hadoop.ozone.container.replication.ReplicationServer.ReplicationConfig; import org.apache.hadoop.util.concurrent.HadoopExecutors; import org.apache.ozone.test.GenericTestUtils; +import org.apache.ozone.test.GenericTestUtils.LogCapturer; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; import org.junit.jupiter.api.io.TempDir; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -149,6 +170,76 @@ public void testStartStopDatanodeStateMachine() throws IOException, } } + @Test + @Timeout(60) + void testDelayedRatisStartupFailureStopsDatanode() throws Exception { + conf.setFromObject(conf.getObject(ReplicationConfig.class).setPort(0)); + DatanodeDetails datanodeDetails = getNewDatanodeDetails(); + ContainerTestUtils.initializeDatanodeLayout(conf, datanodeDetails); + CountDownLatch initializing = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + CountDownLatch shutdown = new CountDownLatch(1); + HddsDatanodeStopService stopService = mock(HddsDatanodeStopService.class); + doAnswer(invocation -> { + shutdown.countDown(); + return null; + }).when(stopService).stopService(); + DatanodeStateMachine stateMachine = new DatanodeStateMachine(null, datanodeDetails, conf, null, null, + stopService, new ReconfigurationHandler("DN", conf, op -> { })); + try (LogCapturer startupLogs = LogCapturer.captureLogs(VersionEndpointTask.class); + LogCapturer stateLogs = LogCapturer.captureLogs(RunningDatanodeState.class)) { + OzoneContainer container = stateMachine.getContainer(); + XceiverServerSpi writeChannel = spy(container.getWriteChannel()); + Field writeChannelField = OzoneContainer.class.getDeclaredField("writeChannel"); + writeChannelField.setAccessible(true); + writeChannelField.set(container, writeChannel); + IllegalStateException failure = new IllegalStateException("Failed to initRaftLog", + new IOException("Corrupt Raft log")); + doAnswer(invocation -> { + initializing.countDown(); + assertThat(release.await(30, TimeUnit.SECONDS)).isTrue(); + throw failure; + }).when(writeChannel).start(); + + stateMachine.startDaemon(); + assertThat(initializing.await(10, TimeUnit.SECONDS)).isTrue(); + // Replication is already serving when Ratis startup fails. + try (Socket socket = new Socket("127.0.0.1", container.getReplicationServer().getPort())) { + assertThat(socket.isConnected()).isTrue(); + } + String clusterId = mockServers.get(0).getClusterId(); + CompletableFuture waitingThread = new CompletableFuture<>(); + Future waitingCaller = executorService.submit(() -> { + waitingThread.complete(Thread.currentThread()); + container.start(clusterId); + return null; + }); + Thread caller = waitingThread.get(10, TimeUnit.SECONDS); + GenericTestUtils.waitFor(() -> caller.getState() == Thread.State.BLOCKED, 10, 10000); + // Let the real heartbeat wait expire before completing the initialization attempt. + GenericTestUtils.waitFor(() -> stateLogs.getOutput().contains("Detected timeout"), 10, 10000); + assertThat(stateMachine.getContext().getState()).isEqualTo(DatanodeStates.RUNNING); + assertThat(shutdown.getCount()).isEqualTo(1); + release.countDown(); + + assertThat(assertThrows(ExecutionException.class, () -> waitingCaller.get(10, TimeUnit.SECONDS)).getCause()) + .isInstanceOf(IOException.class).hasCause(failure); + assertThat(shutdown.await(10, TimeUnit.SECONDS)).isTrue(); + stateMachine.getStateMachineThread().join(10000); + assertThat(stateMachine.getStateMachineThread().isAlive()).isFalse(); + assertThat(stateMachine.getContext().getState()).isEqualTo(DatanodeStates.SHUTDOWN); + assertThat(stateMachine.getContext().getShutdownOnError()).isTrue(); + assertThat(startupLogs.getOutput()).contains("Failed to start required container services", + "java.lang.IllegalStateException: Failed to initRaftLog", "Caused by: java.io.IOException: Corrupt Raft log"); + assertThat(assertThrows(IOException.class, () -> container.start(clusterId))).hasCause(failure); + verify(writeChannel, times(1)).start(); + verify(stopService, times(1)).stopService(); + } finally { + release.countDown(); + stateMachine.stopDaemon(); + } + } + /** * This test explores the state machine by invoking each call in sequence just * like as if the state machine would call it. Because this is a test we are diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/states/datanode/TestRunningDatanodeState.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/states/datanode/TestRunningDatanodeState.java index 51390ae632b5..459dc663c27b 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/states/datanode/TestRunningDatanodeState.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/states/datanode/TestRunningDatanodeState.java @@ -19,20 +19,41 @@ import static org.apache.hadoop.ozone.container.common.statemachine.EndpointStateMachine.EndPointStates.SHUTDOWN; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import java.io.IOException; import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; import java.util.List; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorCompletionService; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; +import org.apache.hadoop.hdds.conf.OzoneConfiguration; +import org.apache.hadoop.hdds.scm.net.HostAndPort; +import org.apache.hadoop.ozone.OzoneConsts; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine.DatanodeStates; import org.apache.hadoop.ozone.container.common.statemachine.EndpointStateMachine; +import org.apache.hadoop.ozone.container.common.statemachine.EndpointStateMachine.EndPointStates; import org.apache.hadoop.ozone.container.common.statemachine.SCMConnectionManager; +import org.apache.hadoop.ozone.container.common.statemachine.StateContext; +import org.apache.hadoop.ozone.container.common.states.endpoint.VersionEndpointTask; +import org.apache.hadoop.ozone.container.ozoneimpl.OzoneContainer; +import org.apache.hadoop.ozone.protocol.VersionResponse; +import org.apache.hadoop.ozone.protocolPB.StorageContainerDatanodeProtocolClientSideTranslatorPB; import org.apache.hadoop.util.Time; +import org.apache.ozone.test.GenericTestUtils.LogCapturer; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; /** * Test class for RunningDatanodeState. @@ -71,20 +92,110 @@ public void testAwait() throws InterruptedException { long endTime = Time.monotonicNow(); assertThat(endTime - startTime).isGreaterThanOrEqualTo(500); + // The next heartbeat must still collect old completions, even if endpoints were removed. + state.clear(); + stateMachines.clear(); + state.execute(executorService); futureOne.complete(SHUTDOWN); - CompletableFuture futureTwo = - new CompletableFuture<>(); - for (int i = 0; i < threadPoolSize; i++) { - ecs.submit(() -> futureTwo.get()); - } - futureTwo.complete(SHUTDOWN); - startTime = Time.monotonicNow(); - state.await(500, TimeUnit.MILLISECONDS); + assertThat(state.await(500, TimeUnit.MILLISECONDS)).isEqualTo(DatanodeStates.SHUTDOWN); endTime = Time.monotonicNow(); assertThat(endTime - startTime).isLessThan(500); executorService.shutdown(); } + + @ParameterizedTest + @ValueSource(ints = {0, 1, 2, 3}) + void testStartupCompletesAfterHeartbeatTimeout(int failureType) throws Exception { + OzoneConfiguration conf = new OzoneConfiguration(); + OzoneContainer container = mock(OzoneContainer.class); + DatanodeStateMachine datanode = mock(DatanodeStateMachine.class); + SCMConnectionManager connections = mock(SCMConnectionManager.class); + when(datanode.getConnectionManager()).thenReturn(connections); + when(datanode.getContainer()).thenReturn(container); + StateContext context = new StateContext(conf, DatanodeStates.RUNNING, datanode, ""); + StorageContainerDatanodeProtocolClientSideTranslatorPB scm = + mock(StorageContainerDatanodeProtocolClientSideTranslatorPB.class); + when(scm.getVersion(null)).thenReturn(versionResponse().getProtobufMessage()); + CountDownLatch initializing = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + Throwable failure = failureType == 1 ? new IOException("Corrupt Raft log") + : failureType == 2 ? new IllegalStateException("Failed to initRaftLog", new IOException("Corrupt Raft log")) + : new AssertionError("Failed to initRaftLog", new IOException("Corrupt Raft log")); + doAnswer(invocation -> { + initializing.countDown(); + assertThat(release.await(10, TimeUnit.SECONDS)).isTrue(); + if (failureType != 0) { + throw failure; + } + return null; + }).when(container).start("cluster"); + ExecutorService executor = Executors.newCachedThreadPool(); + try (LogCapturer logs = LogCapturer.captureLogs(VersionEndpointTask.class); + EndpointStateMachine endpoint = new EndpointStateMachine(new HostAndPort("scm", 9861), scm, conf, "")) { + when(connections.getValues()).thenReturn(Collections.singletonList(endpoint)); + // Keep startup blocked until both the inner future wait and this heartbeat have finished. + context.execute(executor, 5, TimeUnit.SECONDS); + assertThat(initializing.getCount()).isZero(); + assertThat(context.getState()).isEqualTo(DatanodeStates.RUNNING); + assertThat(endpoint.getState()).isEqualTo(EndPointStates.GETVERSION); + release.countDown(); + // The endpoint executor is serial; this barrier waits for the startup task to finish. + endpoint.getExecutorService().submit(() -> { }).get(5, TimeUnit.SECONDS); + if (failureType == 0) { + assertThat(endpoint.getState()).isEqualTo(EndPointStates.REGISTER); + assertThat(context.getShutdownOnError()).isFalse(); + } else { + assertThat(endpoint.getState()).isEqualTo(SHUTDOWN); + assertThat(logs.getOutput()).contains("Failed to start required container services", "scm:9861", + failure.getClass().getName(), "Corrupt Raft log"); + context.execute(executor, 5, TimeUnit.SECONDS); + assertThat(context.getState()).isEqualTo(DatanodeStates.SHUTDOWN); + assertThat(context.getShutdownOnError()).isTrue(); + } + verify(container, times(1)).start("cluster"); + } finally { + release.countDown(); + executor.shutdownNow(); + } + } + + @Test + void testUnavailableScmDoesNotPreventStartup() throws Exception { + OzoneConfiguration conf = new OzoneConfiguration(); + OzoneContainer container = mock(OzoneContainer.class); + DatanodeStateMachine datanode = mock(DatanodeStateMachine.class); + SCMConnectionManager connections = mock(SCMConnectionManager.class); + when(datanode.getConnectionManager()).thenReturn(connections); + when(datanode.getContainer()).thenReturn(container); + StateContext context = new StateContext(conf, DatanodeStates.RUNNING, datanode, ""); + StorageContainerDatanodeProtocolClientSideTranslatorPB unavailable = + mock(StorageContainerDatanodeProtocolClientSideTranslatorPB.class); + when(unavailable.getVersion(null)).thenThrow(new IOException("SCM unavailable")); + StorageContainerDatanodeProtocolClientSideTranslatorPB available = + mock(StorageContainerDatanodeProtocolClientSideTranslatorPB.class); + when(available.getVersion(null)).thenReturn(versionResponse().getProtobufMessage()); + ExecutorService executor = Executors.newCachedThreadPool(); + try (EndpointStateMachine first = new EndpointStateMachine(new HostAndPort("scm1", 9861), unavailable, conf, ""); + EndpointStateMachine second = new EndpointStateMachine(new HostAndPort("scm2", 9861), available, conf, "")) { + when(connections.getValues()).thenReturn(Arrays.asList(first, second)); + context.execute(executor, 5, TimeUnit.SECONDS); + assertThat(first.getState()).isEqualTo(EndPointStates.GETVERSION); + assertThat(second.getState()).isEqualTo(EndPointStates.REGISTER); + assertThat(context.getState()).isEqualTo(DatanodeStates.RUNNING); + assertThat(context.getShutdownOnError()).isFalse(); + verify(container, times(1)).start("cluster"); + } finally { + executor.shutdownNow(); + } + } + + private static VersionResponse versionResponse() { + return VersionResponse.newBuilder().setVersion(1) + .addValue(OzoneConsts.SCM_ID, "scm") + .addValue(OzoneConsts.CLUSTER_ID, "cluster") + .build(); + } } diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/ozoneimpl/TestOzoneContainer.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/ozoneimpl/TestOzoneContainer.java index 91c3f8ed58c2..9f8f9fe4604f 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/ozoneimpl/TestOzoneContainer.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/ozoneimpl/TestOzoneContainer.java @@ -19,15 +19,21 @@ import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result.DISK_OUT_OF_SPACE; import static org.apache.hadoop.ozone.container.common.ContainerTestUtils.createDbInstancesForTestIfNeeded; +import static org.assertj.core.api.Assertions.assertThat; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import com.google.common.base.Preconditions; +import com.google.common.util.concurrent.ThreadFactoryBuilder; import java.io.File; +import java.io.IOException; import java.nio.file.Files; import java.nio.file.Path; import java.util.ArrayList; @@ -39,6 +45,12 @@ import java.util.Set; import java.util.UUID; import java.util.concurrent.ConcurrentSkipListSet; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; import org.apache.commons.io.FileUtils; import org.apache.hadoop.conf.StorageUnit; import org.apache.hadoop.hdds.HddsConfigKeys; @@ -51,11 +63,15 @@ import org.apache.hadoop.hdds.utils.db.Table; import org.apache.hadoop.ozone.OzoneConfigKeys; import org.apache.hadoop.ozone.container.common.ContainerTestUtils; +import org.apache.hadoop.ozone.container.common.SCMTestUtils; import org.apache.hadoop.ozone.container.common.helpers.BlockData; import org.apache.hadoop.ozone.container.common.helpers.ChunkInfo; import org.apache.hadoop.ozone.container.common.impl.ContainerLayoutVersion; import org.apache.hadoop.ozone.container.common.impl.ContainerSet; import org.apache.hadoop.ozone.container.common.interfaces.DBHandle; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine.DatanodeStates; +import org.apache.hadoop.ozone.container.common.statemachine.StateContext; import org.apache.hadoop.ozone.container.common.utils.StorageVolumeUtil; import org.apache.hadoop.ozone.container.common.volume.HddsVolume; import org.apache.hadoop.ozone.container.common.volume.MutableVolumeSet; @@ -67,6 +83,8 @@ import org.apache.hadoop.ozone.container.keyvalue.helpers.BlockUtils; import org.junit.jupiter.api.AfterEach; import org.junit.jupiter.api.io.TempDir; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; /** * This class is used to test OzoneContainer. @@ -116,6 +134,69 @@ public void cleanUp() { } } + @ParameterizedTest + @ValueSource(ints = {0, 1, 2, 3}) + void testConcurrentStartup(int failureType) throws Exception { + conf = SCMTestUtils.getConf(folder.toFile()); + conf.setBoolean(OzoneConfigKeys.HDDS_CONTAINER_RATIS_IPC_RANDOM_PORT, true); + conf.setBoolean(OzoneConfigKeys.HDDS_CONTAINER_IPC_RANDOM_PORT, true); + ContainerTestUtils.initializeDatanodeLayout(conf, datanodeDetails); + OzoneContainer container = spy(ContainerTestUtils.getOzoneContainer(datanodeDetails, conf)); + CountDownLatch initializing = new CountDownLatch(1); + CountDownLatch release = new CountDownLatch(1); + CountDownLatch secondCaller = new CountDownLatch(1); + Throwable failure = failureType == 1 ? new IOException("Cannot read container metadata") + : failureType == 2 ? new IllegalStateException("Failed to initRaftLog", new IOException("Corrupt Raft log")) + : new AssertionError("Failed to initRaftLog", new IOException("Corrupt Raft log")); + ExecutorService executor = Executors.newFixedThreadPool(2, + new ThreadFactoryBuilder().setDaemon(true).build()); + try { + doAnswer(invocation -> { + initializing.countDown(); + assertThat(release.await(10, TimeUnit.SECONDS)).isTrue(); + if (failureType != 0) { + throw failure; + } + return invocation.callRealMethod(); + }).when(container).buildContainerSet(); + Future first = executor.submit(() -> { + container.start(clusterId); + return null; + }); + assertThat(initializing.await(10, TimeUnit.SECONDS)).isTrue(); + // Full reports use the container monitor and must not wait for local service startup. + DatanodeStateMachine datanode = mock(DatanodeStateMachine.class); + when(datanode.getContainer()).thenReturn(container); + StateContext reportContext = new StateContext(conf, DatanodeStates.RUNNING, datanode, ""); + assertThat(executor.submit(reportContext::getFullContainerReportDiscardPendingICR) + .get(5, TimeUnit.SECONDS).getReportsCount()).isZero(); + Future second = executor.submit(() -> { + secondCaller.countDown(); + container.start(clusterId); + return null; + }); + assertThat(secondCaller.await(10, TimeUnit.SECONDS)).isTrue(); + assertThat(second.isDone()).isFalse(); + release.countDown(); + if (failureType == 0) { + first.get(10, TimeUnit.SECONDS); + second.get(10, TimeUnit.SECONDS); + container.start(clusterId); + } else { + assertThat(assertThrows(ExecutionException.class, () -> first.get(10, TimeUnit.SECONDS))) + .hasCause(failure); + assertThat(assertThrows(ExecutionException.class, () -> second.get(10, TimeUnit.SECONDS)).getCause()) + .isInstanceOf(IOException.class).hasCause(failure); + assertThat(assertThrows(IOException.class, () -> container.start(clusterId))).hasCause(failure); + } + verify(container, times(1)).buildContainerSet(); + } finally { + release.countDown(); + executor.shutdownNow(); + container.stop(); + } + } + /** * Create a mock {@link HddsVolume} to track container IDs. */