diff --git a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java index e75a6ca85d40c..dcf73ce3d4a29 100644 --- a/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java +++ b/core/src/test/java/kafka/log/remote/RemoteLogManagerTest.java @@ -126,6 +126,7 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; import java.util.function.BiConsumer; import java.util.function.Supplier; @@ -1717,9 +1718,29 @@ void testFetchOffsetByTimestampWithTieredStorageDoesNotFetchIndexWhenExistsLocal FileRecords.TimestampAndOffset expectedRemoteResult = new FileRecords.TimestampAndOffset(timestamp + 999, 999, Optional.of(Integer.MAX_VALUE)); Partition mockFollowerPartition = mockPartition(tpId); - LogSegment logSegment = mockLogSegment(50L, timestamp, null); - LogSegment logSegment1 = mockLogSegment(100L, timestamp + 1, expectedLocalResult); - when(mockLog.logSegments()).thenReturn(Arrays.asList(logSegment, logSegment1)); + LogSegment logSegmentBaseOffset50 = mockLogSegment(50L, timestamp, null); + LogSegment logSegmentBaseOffset100 = mockLogSegment(100L, timestamp + 1, expectedLocalResult); + LogSegment logSegmentBaseOffset101 = mockLogSegment(101L, timestamp + 1, expectedLocalResult); + + // Constants representing the states of local log segments + final int twoSegmentsBaseOffsets50and100 = 0; + final int oneSegmentBaseOffset100 = 1; + final int oneSegmentBaseOffset101 = 2; + + AtomicInteger localLogOffsetState = new AtomicInteger(twoSegmentsBaseOffsets50and100); + + when(mockLog.logSegments()).thenAnswer(invocation -> { + if (localLogOffsetState.get() == twoSegmentsBaseOffsets50and100) { + return Arrays.asList(logSegmentBaseOffset50, logSegmentBaseOffset100); + } else if (localLogOffsetState.get() == oneSegmentBaseOffset100) { + return Collections.singletonList(logSegmentBaseOffset100); + } else if (localLogOffsetState.get() == oneSegmentBaseOffset101) { + return Collections.singletonList(logSegmentBaseOffset101); + } else { + throw new IllegalStateException("Unexpected localLogOffsetState"); + } + }); + when(mockLog.logEndOffset()).thenReturn(300L); remoteLogManager = new RemoteLogManager(config.remoteLogManagerConfig(), brokerId, logDir, clusterId, time, partition -> Optional.of(mockLog), @@ -1745,12 +1766,12 @@ Optional lookupTimestamp(RemoteLogSegmentMetadat // Move the local-log start offset to 100L, still the read from the remote storage should be short-circuited // as the message with (timestamp + 1) exists in the local log - when(mockLog.logSegments()).thenReturn(Collections.singletonList(logSegment1)); + localLogOffsetState.set(oneSegmentBaseOffset100); assertEquals(Optional.of(expectedLocalResult), remoteLogManager.findOffsetByTimestamp(tp, timestamp + 1, 0L, cache)); // Move the local log start offset to 101L, now message with (timestamp + 1) does not exist in the local log and // the indexes needs to be fetched from the remote storage - when(logSegment1.baseOffset()).thenReturn(101L); + localLogOffsetState.set(oneSegmentBaseOffset101); assertEquals(Optional.of(expectedRemoteResult), remoteLogManager.findOffsetByTimestamp(tp, timestamp + 1, 0L, cache)); }