diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/TestOzoneFileSystemDataStreamEnablement.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/TestOzoneFileSystemDataStreamEnablement.java index 11a55b471bd..73fe728db6e 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/TestOzoneFileSystemDataStreamEnablement.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/TestOzoneFileSystemDataStreamEnablement.java @@ -214,6 +214,11 @@ private List openRatisThreePipelines() { Pipeline.PipelineState.OPEN); } + private static boolean noNodesHaveDatastreamPort(Pipeline pipeline) { + return pipeline.getNodes().stream() + .noneMatch(n -> n.hasPort(RATIS_DATASTREAM)); + } + private static boolean allNodesHaveDatastreamPort(Pipeline pipeline) { return pipeline.getNodes().stream() .allMatch(n -> n.hasPort(RATIS_DATASTREAM)); @@ -336,7 +341,7 @@ public void testCloseNonStreamablePipelineThenStream() throws Exception { } final List before = openRatisThreePipelines(); assertFalse(before.isEmpty()); - before.forEach(p -> assertFalse(allNodesHaveDatastreamPort(p), + before.forEach(p -> assertTrue(noNodesHaveDatastreamPort(p), "pipeline should be portless before enabling datastream")); rollingRestartEnablingDataStream(); @@ -350,8 +355,8 @@ public void testCloseNonStreamablePipelineThenStream() throws Exception { final PipelineManagerImpl pipelineManager = (PipelineManagerImpl) cluster.getStorageContainerManager().getPipelineManager(); final List reloaded = openRatisThreePipelines(); - assertFalse(reloaded.isEmpty()); - reloaded.forEach(p -> assertFalse(allNodesHaveDatastreamPort(p), + reloaded.retainAll(before); + reloaded.forEach(p -> assertTrue(noNodesHaveDatastreamPort(p), "reloaded pipeline should still be portless")); // Close the pipeline(s) exposing the new datastream port; a fresh