From d51bf48efded408219d11c1fb15545cf904841e3 Mon Sep 17 00:00:00 2001 From: Shourya Dutta Biswas <114977491+shourya035@users.noreply.github.com> Date: Fri, 21 Aug 2026 15:55:23 +0530 Subject: [PATCH] Verify remote translog size-based flushing Cover legacy moving-boundary behavior with byte tracking disabled. Verify accumulated bytes trigger a commit when tracking is enabled. Signed-off-by: Shourya Dutta Biswas <114977491+shourya035@users.noreply.github.com> --- .../remotestore/RemoteStoreCoreTestCase.java | 105 ++++++++++++++++++ 1 file changed, 105 insertions(+) diff --git a/test/framework/src/main/java/org/opensearch/remotestore/RemoteStoreCoreTestCase.java b/test/framework/src/main/java/org/opensearch/remotestore/RemoteStoreCoreTestCase.java index d34db204a112f..4c4c082ed1a05 100644 --- a/test/framework/src/main/java/org/opensearch/remotestore/RemoteStoreCoreTestCase.java +++ b/test/framework/src/main/java/org/opensearch/remotestore/RemoteStoreCoreTestCase.java @@ -18,6 +18,7 @@ import org.opensearch.action.admin.indices.settings.put.UpdateSettingsRequest; import org.opensearch.action.index.IndexResponse; import org.opensearch.action.search.SearchPhaseExecutionException; +import org.opensearch.action.support.WriteRequest; import org.opensearch.cluster.health.ClusterHealthStatus; import org.opensearch.cluster.metadata.IndexMetadata; import org.opensearch.cluster.routing.RecoverySource; @@ -28,10 +29,13 @@ import org.opensearch.common.unit.TimeValue; import org.opensearch.common.util.concurrent.BufferedAsyncIOProcessor; import org.opensearch.index.IndexSettings; +import org.opensearch.index.seqno.SequenceNumbers; import org.opensearch.index.shard.IndexShard; import org.opensearch.index.shard.IndexShardClosedException; +import org.opensearch.index.translog.RemoteFsTranslog; import org.opensearch.index.translog.Translog; import org.opensearch.index.translog.Translog.Durability; +import org.opensearch.indices.IndexingMemoryController; import org.opensearch.indices.IndicesService; import org.opensearch.indices.RemoteStoreSettings; import org.opensearch.indices.recovery.PeerRecoveryTargetService; @@ -75,6 +79,7 @@ import static org.opensearch.test.hamcrest.OpenSearchAssertions.assertHitCount; import static org.hamcrest.Matchers.comparesEqualTo; import static org.hamcrest.Matchers.greaterThan; +import static org.hamcrest.Matchers.greaterThanOrEqualTo; import static org.hamcrest.Matchers.is; import static org.hamcrest.Matchers.oneOf; @@ -995,6 +1000,106 @@ public void testFlushOnTooManyRemoteTranslogFiles() throws Exception { } } + /** + * Verifies that disabling byte tracking preserves the legacy remote-translog flush behavior. + *
    + *
  1. Disables reader-count and background flush triggers so only the size threshold can initiate a periodic flush.
  2. + *
  3. Creates a remote-backed index with byte tracking disabled and a small translog flush threshold.
  4. + *
  5. Indexes and refreshes documents until the remote retention boundary advances beyond its initial value.
  6. + *
  7. Asserts that the moving legacy boundary prevents a periodic flush and leaves the commit checkpoint unchanged.
  8. + *
+ */ + public void testRemoteTranslogSizeFlushUsesLegacyRetentionBoundaryWhenDisabled() throws Exception { + assertRemoteTranslogSizeBasedFlush(false); + } + + /** + * Verifies that enabled byte tracking restores size-based remote-translog flushes across refreshes. + *
    + *
  1. Disables reader-count and background flush triggers so only the size threshold can initiate a periodic flush.
  2. + *
  3. Creates a remote-backed index with byte tracking enabled and a small translog flush threshold.
  4. + *
  5. Indexes and refreshes documents while the remote retention boundary advances and tracked bytes accumulate.
  6. + *
  7. Asserts that a periodic flush commits the accumulated operations after the threshold is crossed.
  8. + *
  9. Asserts that the successful commit resets the size trigger.
  10. + *
+ */ + public void testRemoteTranslogSizeFlushTracksBytesAcrossRefreshes() throws Exception { + assertRemoteTranslogSizeBasedFlush(true); + } + + private void assertRemoteTranslogSizeBasedFlush(boolean trackBytesSinceLastCommit) throws Exception { + Settings nodeSettings = Settings.builder().put(IndexingMemoryController.SHARD_INACTIVE_TIME_SETTING.getKey(), "1h").build(); + internalCluster().startClusterManagerOnlyNode(nodeSettings); + String dataNode = internalCluster().startDataOnlyNode(nodeSettings); + + assertAcked( + client().admin() + .cluster() + .prepareUpdateSettings() + .setPersistentSettings( + Settings.builder() + .put(RemoteStoreSettings.CLUSTER_REMOTE_MAX_TRANSLOG_READERS.getKey(), -1) + .put(CLUSTER_REMOTE_TRANSLOG_BUFFER_INTERVAL_SETTING.getKey(), "0ms") + .put( + RemoteStoreSettings.CLUSTER_REMOTE_TRANSLOG_TRACK_BYTES_SINCE_LAST_COMMIT_SETTING.getKey(), + trackBytesSinceLastCommit + ) + ) + .get() + ); + + createIndex( + INDEX_NAME, + Settings.builder() + .put(remoteStoreIndexSettings(0, 10000L, -1)) + .put(IndexSettings.INDEX_TRANSLOG_FLUSH_THRESHOLD_SIZE_SETTING.getKey(), "32kb") + .put(IndexSettings.INDEX_PERIODIC_FLUSH_INTERVAL_SETTING.getKey(), "-1") + .build() + ); + ensureGreen(INDEX_NAME); + + IndexShard indexShard = getIndexShard(dataNode, INDEX_NAME); + RemoteFsTranslog remoteTranslog = (RemoteFsTranslog) getTranslog(indexShard); + long initialRemoteRetentionBoundary = remoteTranslog.getMinUnreferencedSeqNoInSegments(0); + long initialPeriodicFlushes = indexShard.flushStats().getPeriodic(); + long initialCommitCheckpoint = getLastCommittedLocalCheckpoint(indexShard); + assertFalse(indexShard.shouldPeriodicallyFlush()); + + String payload = randomAlphaOfLength(4 * 1024); + for (int i = 0; i < 10; i++) { + IndexResponse response = client(dataNode).prepareIndex(INDEX_NAME) + .setId(Integer.toString(i)) + .setSource("payload", payload) + .setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE) + .get(); + assertBusy( + () -> assertThat(remoteTranslog.getMinUnreferencedSeqNoInSegments(0), greaterThanOrEqualTo(response.getSeqNo())), + 30, + TimeUnit.SECONDS + ); + if (trackBytesSinceLastCommit == false) { + assertFalse(indexShard.shouldPeriodicallyFlush()); + } + } + assertThat(remoteTranslog.getMinUnreferencedSeqNoInSegments(0), greaterThan(initialRemoteRetentionBoundary)); + + if (trackBytesSinceLastCommit) { + assertBusy(() -> { + assertThat(indexShard.flushStats().getPeriodic(), greaterThan(initialPeriodicFlushes)); + assertThat(getLastCommittedLocalCheckpoint(indexShard), greaterThan(initialCommitCheckpoint)); + assertFalse(indexShard.shouldPeriodicallyFlush()); + }, 30, TimeUnit.SECONDS); + } else { + assertEquals(initialPeriodicFlushes, indexShard.flushStats().getPeriodic()); + assertEquals(initialCommitCheckpoint, getLastCommittedLocalCheckpoint(indexShard)); + assertFalse(indexShard.shouldPeriodicallyFlush()); + } + } + + private long getLastCommittedLocalCheckpoint(IndexShard indexShard) { + return Long.parseLong(indexShard.commitStats().getUserData().get(SequenceNumbers.LOCAL_CHECKPOINT_KEY)); + } + public void testAsyncTranslogDurabilityRestrictionsThroughIdxTemplates() throws Exception { logger.info("Starting up cluster manager with cluster.remote_store.index.restrict.async-durability set to true"); String cm1 = internalCluster().startClusterManagerOnlyNode(