ss77892 commented on code in PR #11080: URL: https://github.com/apache/ozone/pull/11080#discussion_r4169844841
########## hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/transport/server/TestGrpcXceiverService.java: ########## @@ -0,0 +1,165 @@ +/* + * 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.ozone.container.common.transport.server; + +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; + +import java.io.File; +import java.nio.file.Files; +import java.nio.file.Path; +import java.time.Duration; +import java.util.concurrent.atomic.AtomicReference; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandResponseProto; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.DatanodeBlockID; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ReadBlockRequestProto; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Type; +import org.apache.hadoop.hdds.utils.io.RandomAccessFileChannel; +import org.apache.hadoop.ozone.container.common.interfaces.ContainerDispatcher; +import org.apache.ozone.test.GenericTestUtils; +import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +/** + * Tests for the streaming ReadBlock handling in {@link GrpcXceiverService}: the block file held open for a + * stream is closed once the stream is idle and reopened by the next request. + */ +class TestGrpcXceiverService { + + private static final Duration IDLE_TIMEOUT = Duration.ofMillis(200); + + @TempDir + private Path tempDir; + + private GrpcXceiverService service; + + @AfterEach + void shutdown() { + if (service != null) { + service.shutdown(); + } + } + + private static ContainerCommandRequestProto readBlockRequest() { + return ContainerCommandRequestProto.newBuilder() + .setCmdType(Type.ReadBlock) + .setContainerID(1) + .setDatanodeUuid("dn") + .setReadBlock(ReadBlockRequestProto.newBuilder() + .setBlockID(DatanodeBlockID.newBuilder().setContainerID(1).setLocalID(1)) + .setOffset(0)) + .build(); + } + + /** + * Mimics {@code KeyValueHandler.readBlockImpl}: open the block file on first use, then serve the request. + * The optional hook runs while the request is in flight. + */ + private ContainerDispatcher mockDispatcher(File blockFile, AtomicReference<RandomAccessFileChannel> channelRef, + Runnable inFlight) throws Exception { + ContainerDispatcher dispatcher = mock(ContainerDispatcher.class); + doAnswer(inv -> { + RandomAccessFileChannel channel = inv.getArgument(2); + channelRef.set(channel); + if (!channel.isOpen()) { + channel.open(blockFile); + } + inFlight.run(); + return null; + }).when(dispatcher).streamDataReadOnly(any(), any(), any(), any()); + return dispatcher; + } + + @Test + void idleStreamClosesBlockFileAndNextRequestReopensIt() throws Exception { + File blockFile = Files.createFile(tempDir.resolve("block")).toFile(); + AtomicReference<RandomAccessFileChannel> channelRef = new AtomicReference<>(); + service = new GrpcXceiverService(mockDispatcher(blockFile, channelRef, () -> { }), IDLE_TIMEOUT, "test-"); + StreamObserver<ContainerCommandResponseProto> responseObserver = mock(StreamObserver.class); + StreamObserver<ContainerCommandRequestProto> requestObserver = service.send(responseObserver); + + requestObserver.onNext(readBlockRequest()); + RandomAccessFileChannel channel = channelRef.get(); + assertNotNull(channel); + assertTrue(channel.isOpen(), "block file should be open right after a request"); + + GenericTestUtils.waitFor(() -> !channel.isOpen(), 20, 5000); + verify(responseObserver, never()).onError(any()); + verify(responseObserver, never()).onCompleted(); + + requestObserver.onNext(readBlockRequest()); + assertTrue(channel.isOpen(), "next request should reopen the block file"); Review Comment: Fixed. The open state is now recorded inside the dispatcher hook, while the request lock is held, so the idle closer can't close the file before the check. The close is still verified with waitFor. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
