From 622595fd55f7b76f48047560da5d5069b6ca4d41 Mon Sep 17 00:00:00 2001 From: gaoran10 Date: Tue, 30 Jun 2020 10:28:54 +0800 Subject: [PATCH 1/7] set the default offload deletion lag value to -1 --- .../org/apache/pulsar/common/policies/data/OffloadPolicies.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java index 4936923dfda6b..a743c84135079 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java @@ -41,7 +41,7 @@ public class OffloadPolicies { public final static String[] DRIVER_NAMES = {"S3", "aws-s3", "google-cloud-storage", "filesystem"}; public final static String DEFAULT_OFFLOADER_DIRECTORY = "./offloaders"; public final static long DEFAULT_OFFLOAD_THRESHOLD_IN_BYTES = -1; - public final static Long DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS = null; + public final static long DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS = -1; // common config private String offloadersDirectory = DEFAULT_OFFLOADER_DIRECTORY; From f683f882114eabeca18bc8d3acb4b3accfe9e227 Mon Sep 17 00:00:00 2001 From: gaoran10 Date: Tue, 30 Jun 2020 11:21:03 +0800 Subject: [PATCH 2/7] add unit test for method `ManagedLedgerImpl.isOffloadedNeedsDelete` --- .../mledger/impl/OffloadLedgerDeleteTest.java | 30 +++++++++++++++++++ 1 file changed, 30 insertions(+) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java index 8fbc588b56086..d0c31a7ad829b 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java @@ -20,15 +20,21 @@ import static org.apache.bookkeeper.mledger.impl.OffloadPrefixTest.assertEventuallyTrue; +import java.lang.reflect.Method; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; +import org.apache.bookkeeper.mledger.LedgerOffloader; import org.apache.bookkeeper.mledger.ManagedCursor; +import org.apache.bookkeeper.mledger.ManagedLedger; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; +import org.apache.bookkeeper.mledger.proto.MLDataFormats; import org.apache.bookkeeper.mledger.util.MockClock; import org.apache.bookkeeper.test.MockedBookKeeperTestCase; +import org.apache.pulsar.common.policies.data.OffloadPolicies; +import org.mockito.Mockito; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -201,4 +207,28 @@ public void testLaggedDeleteSlowConsumer() throws Exception { .map(e -> e.getLedgerId()).collect(Collectors.toSet()), offloader.offloadedLedgers()); } + + @Test + public void isOffloadedNeedsDeleteTest() throws Exception { + LedgerOffloader ledgerOffloader = Mockito.mock(LedgerOffloader.class); + Mockito.when(ledgerOffloader.getOffloadPolicies()).thenReturn(new OffloadPolicies()); + + ManagedLedgerConfig config = new ManagedLedgerConfig(); + MockClock clock = new MockClock(); + config.setLedgerOffloader(ledgerOffloader); + config.setClock(clock); + + ManagedLedger managedLedger = factory.open("isOffloadedNeedsDeleteTest", config); + Class clazz = ManagedLedgerImpl.class; + Method method = clazz.getDeclaredMethod("isOffloadedNeedsDelete", MLDataFormats.OffloadContext.class); + method.setAccessible(true); + + MLDataFormats.OffloadContext offloadContext = MLDataFormats.OffloadContext.newBuilder() + .setTimestamp(System.currentTimeMillis() - 1000) + .setComplete(true) + .setBookkeeperDeleted(false) + .build(); + Boolean needsDelete = (Boolean) method.invoke(managedLedger, offloadContext); + Assert.assertTrue(needsDelete); + } } From 55dda16ea335c4b420e39d1957ad0e3d4c8368e6 Mon Sep 17 00:00:00 2001 From: gaoran10 Date: Tue, 30 Jun 2020 11:35:45 +0800 Subject: [PATCH 3/7] fix test --- .../mledger/impl/OffloadLedgerDeleteTest.java | 30 +++++++++++++++++-- 1 file changed, 28 insertions(+), 2 deletions(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java index d0c31a7ad829b..0c43b72e1ebe1 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java @@ -210,8 +210,9 @@ public void testLaggedDeleteSlowConsumer() throws Exception { @Test public void isOffloadedNeedsDeleteTest() throws Exception { + OffloadPolicies offloadPolicies = new OffloadPolicies(); LedgerOffloader ledgerOffloader = Mockito.mock(LedgerOffloader.class); - Mockito.when(ledgerOffloader.getOffloadPolicies()).thenReturn(new OffloadPolicies()); + Mockito.when(ledgerOffloader.getOffloadPolicies()).thenReturn(offloadPolicies); ManagedLedgerConfig config = new ManagedLedgerConfig(); MockClock clock = new MockClock(); @@ -224,11 +225,36 @@ public void isOffloadedNeedsDeleteTest() throws Exception { method.setAccessible(true); MLDataFormats.OffloadContext offloadContext = MLDataFormats.OffloadContext.newBuilder() - .setTimestamp(System.currentTimeMillis() - 1000) + .setTimestamp(config.getClock().millis() - 1000) .setComplete(true) .setBookkeeperDeleted(false) .build(); Boolean needsDelete = (Boolean) method.invoke(managedLedger, offloadContext); Assert.assertTrue(needsDelete); + + offloadContext = MLDataFormats.OffloadContext.newBuilder() + .setTimestamp(config.getClock().millis() - 1000) + .setComplete(false) + .setBookkeeperDeleted(false) + .build(); + needsDelete = (Boolean) method.invoke(managedLedger, offloadContext); + Assert.assertFalse(needsDelete); + + offloadContext = MLDataFormats.OffloadContext.newBuilder() + .setTimestamp(config.getClock().millis() - 1000) + .setComplete(true) + .setBookkeeperDeleted(true) + .build(); + needsDelete = (Boolean) method.invoke(managedLedger, offloadContext); + Assert.assertFalse(needsDelete); + + offloadPolicies.setManagedLedgerOffloadDeletionLagInMillis(1000L * 2); + offloadContext = MLDataFormats.OffloadContext.newBuilder() + .setTimestamp(config.getClock().millis() - 1000) + .setComplete(true) + .setBookkeeperDeleted(false) + .build(); + needsDelete = (Boolean) method.invoke(managedLedger, offloadContext); + Assert.assertFalse(needsDelete); } } From 45f44113a7c14126831a9eb38a5e6a1c9cd275d2 Mon Sep 17 00:00:00 2001 From: gaoran10 Date: Wed, 1 Jul 2020 00:47:00 +0800 Subject: [PATCH 4/7] fix test --- .../bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java | 8 ++++---- .../apache/bookkeeper/mledger/impl/OffloadPrefixTest.java | 2 +- .../apache/pulsar/broker/admin/impl/NamespacesBase.java | 5 ++--- .../pulsar/common/policies/data/OffloadPolicies.java | 2 +- 4 files changed, 8 insertions(+), 9 deletions(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java index 0c43b72e1ebe1..45636fba6f7dc 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java @@ -54,7 +54,7 @@ public void testLaggedDelete() throws Exception { config.setMinimumRolloverTime(0, TimeUnit.SECONDS); config.setRetentionTime(10, TimeUnit.MINUTES); config.setRetentionSizeInMB(10); - offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(new Long(300000)); + offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(300000L); config.setLedgerOffloader(offloader); config.setClock(clock); @@ -116,7 +116,7 @@ public void testLaggedDeleteRetentionSetLower() throws Exception { config.setMinimumRolloverTime(0, TimeUnit.SECONDS); config.setRetentionTime(5, TimeUnit.MINUTES); config.setRetentionSizeInMB(10); - offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(new Long(600000)); + offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(600000L); config.setLedgerOffloader(offloader); config.setClock(clock); @@ -163,7 +163,7 @@ public void testLaggedDeleteSlowConsumer() throws Exception { config.setMaxEntriesPerLedger(10); config.setMinimumRolloverTime(0, TimeUnit.SECONDS); config.setRetentionTime(10, TimeUnit.MINUTES); - offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(new Long(300000)); + offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(300000L); config.setLedgerOffloader(offloader); config.setClock(clock); @@ -248,7 +248,7 @@ public void isOffloadedNeedsDeleteTest() throws Exception { needsDelete = (Boolean) method.invoke(managedLedger, offloadContext); Assert.assertFalse(needsDelete); - offloadPolicies.setManagedLedgerOffloadDeletionLagInMillis(1000L * 2); + offloadPolicies.setManagedLedgerOffloadDeletionLagInMillis(1000 * 2); offloadContext = MLDataFormats.OffloadContext.newBuilder() .setTimestamp(config.getClock().millis() - 1000) .setComplete(true) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java index 97bab568e4cc7..87de7b56f0592 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java @@ -608,7 +608,7 @@ public void testOffloadDelete() throws Exception { config.setMaxEntriesPerLedger(10); config.setMinimumRolloverTime(0, TimeUnit.SECONDS); config.setRetentionTime(0, TimeUnit.MINUTES); - offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(new Long(100)); + offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(100); offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInBytes(100); config.setLedgerOffloader(offloader); ManagedLedgerImpl ledger = (ManagedLedgerImpl)factory.open("my_test_ledger", config); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java index 00f1f4b4f7534..c9d36ba99ebaf 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java @@ -2652,9 +2652,8 @@ protected void internalSetOffloadPolicies(AsyncResponse asyncResponse, OffloadPo final String path = path(POLICIES, namespaceName.toString()); byte[] content = globalZk().getData(path, null, nodeStat); Policies policies = jsonMapper().readValue(content, Policies.class); - if (offloadPolicies.getManagedLedgerOffloadDeletionLagInMillis() == null && OffloadPolicies.DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS == null - || offloadPolicies.getManagedLedgerOffloadDeletionLagInMillis() != null - && offloadPolicies.getManagedLedgerOffloadDeletionLagInMillis().equals(OffloadPolicies.DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS)) { + if (offloadPolicies.getManagedLedgerOffloadDeletionLagInMillis() == + OffloadPolicies.DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS) { offloadPolicies.setManagedLedgerOffloadDeletionLagInMillis(policies.offload_deletion_lag_ms); } else { policies.offload_deletion_lag_ms = offloadPolicies.getManagedLedgerOffloadDeletionLagInMillis(); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java index a743c84135079..80719eec14a6e 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java @@ -49,7 +49,7 @@ public class OffloadPolicies { private int managedLedgerOffloadMaxThreads = DEFAULT_OFFLOAD_MAX_THREADS; private int managedLedgerOffloadPrefetchRounds = DEFAULT_OFFLOAD_MAX_PREFETCH_ROUNDS; private long managedLedgerOffloadThresholdInBytes = DEFAULT_OFFLOAD_THRESHOLD_IN_BYTES; - private Long managedLedgerOffloadDeletionLagInMillis = DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS; + private long managedLedgerOffloadDeletionLagInMillis = DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS; // s3 config, set by service configuration or cli private String s3ManagedLedgerOffloadRegion = null; From 7ae01ef4e22770e50317e1c304cfc25e251f4350 Mon Sep 17 00:00:00 2001 From: gaoran10 Date: Sun, 5 Jul 2020 17:43:59 +0800 Subject: [PATCH 5/7] fix --- .../apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 4 +++- .../bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java | 2 +- .../org/apache/pulsar/broker/admin/impl/NamespacesBase.java | 5 +++-- .../apache/pulsar/common/policies/data/OffloadPolicies.java | 4 ++-- 4 files changed, 9 insertions(+), 6 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java index f45b341d61d86..ea0614d478e93 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java @@ -1960,7 +1960,9 @@ private boolean isOffloadedNeedsDelete(OffloadContext offload) { long elapsedMs = clock.millis() - offload.getTimestamp(); if (config.getLedgerOffloader() != null && config.getLedgerOffloader() != NullLedgerOffloader.INSTANCE - && config.getLedgerOffloader().getOffloadPolicies() != null) { + && config.getLedgerOffloader().getOffloadPolicies() != null + && config.getLedgerOffloader().getOffloadPolicies() + .getManagedLedgerOffloadDeletionLagInMillis() != null) { return offload.getComplete() && !offload.getBookkeeperDeleted() && elapsedMs > config.getLedgerOffloader() .getOffloadPolicies().getManagedLedgerOffloadDeletionLagInMillis(); diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java index 45636fba6f7dc..5304bc695b126 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java @@ -248,7 +248,7 @@ public void isOffloadedNeedsDeleteTest() throws Exception { needsDelete = (Boolean) method.invoke(managedLedger, offloadContext); Assert.assertFalse(needsDelete); - offloadPolicies.setManagedLedgerOffloadDeletionLagInMillis(1000 * 2); + offloadPolicies.setManagedLedgerOffloadDeletionLagInMillis(1000L * 2); offloadContext = MLDataFormats.OffloadContext.newBuilder() .setTimestamp(config.getClock().millis() - 1000) .setComplete(true) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java index c9d36ba99ebaf..00f1f4b4f7534 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/impl/NamespacesBase.java @@ -2652,8 +2652,9 @@ protected void internalSetOffloadPolicies(AsyncResponse asyncResponse, OffloadPo final String path = path(POLICIES, namespaceName.toString()); byte[] content = globalZk().getData(path, null, nodeStat); Policies policies = jsonMapper().readValue(content, Policies.class); - if (offloadPolicies.getManagedLedgerOffloadDeletionLagInMillis() == - OffloadPolicies.DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS) { + if (offloadPolicies.getManagedLedgerOffloadDeletionLagInMillis() == null && OffloadPolicies.DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS == null + || offloadPolicies.getManagedLedgerOffloadDeletionLagInMillis() != null + && offloadPolicies.getManagedLedgerOffloadDeletionLagInMillis().equals(OffloadPolicies.DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS)) { offloadPolicies.setManagedLedgerOffloadDeletionLagInMillis(policies.offload_deletion_lag_ms); } else { policies.offload_deletion_lag_ms = offloadPolicies.getManagedLedgerOffloadDeletionLagInMillis(); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java index 80719eec14a6e..4936923dfda6b 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/OffloadPolicies.java @@ -41,7 +41,7 @@ public class OffloadPolicies { public final static String[] DRIVER_NAMES = {"S3", "aws-s3", "google-cloud-storage", "filesystem"}; public final static String DEFAULT_OFFLOADER_DIRECTORY = "./offloaders"; public final static long DEFAULT_OFFLOAD_THRESHOLD_IN_BYTES = -1; - public final static long DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS = -1; + public final static Long DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS = null; // common config private String offloadersDirectory = DEFAULT_OFFLOADER_DIRECTORY; @@ -49,7 +49,7 @@ public class OffloadPolicies { private int managedLedgerOffloadMaxThreads = DEFAULT_OFFLOAD_MAX_THREADS; private int managedLedgerOffloadPrefetchRounds = DEFAULT_OFFLOAD_MAX_PREFETCH_ROUNDS; private long managedLedgerOffloadThresholdInBytes = DEFAULT_OFFLOAD_THRESHOLD_IN_BYTES; - private long managedLedgerOffloadDeletionLagInMillis = DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS; + private Long managedLedgerOffloadDeletionLagInMillis = DEFAULT_OFFLOAD_DELETION_LAG_IN_MILLIS; // s3 config, set by service configuration or cli private String s3ManagedLedgerOffloadRegion = null; From 75e02953491bf715f6299cb5fbed449d11d46a9f Mon Sep 17 00:00:00 2001 From: gaoran10 Date: Sun, 5 Jul 2020 18:04:11 +0800 Subject: [PATCH 6/7] fix --- .../org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java index 87de7b56f0592..d2d41489e35b0 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadPrefixTest.java @@ -608,7 +608,7 @@ public void testOffloadDelete() throws Exception { config.setMaxEntriesPerLedger(10); config.setMinimumRolloverTime(0, TimeUnit.SECONDS); config.setRetentionTime(0, TimeUnit.MINUTES); - offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(100); + offloader.getOffloadPolicies().setManagedLedgerOffloadDeletionLagInMillis(100L); offloader.getOffloadPolicies().setManagedLedgerOffloadThresholdInBytes(100); config.setLedgerOffloader(offloader); ManagedLedgerImpl ledger = (ManagedLedgerImpl)factory.open("my_test_ledger", config); From 89a44c964896888943e83c507112eb658dbb3216 Mon Sep 17 00:00:00 2001 From: gaoran10 Date: Mon, 6 Jul 2020 10:57:28 +0800 Subject: [PATCH 7/7] fix test --- .../mledger/impl/OffloadLedgerDeleteTest.java | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java index 5304bc695b126..736cd8bd8362b 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/OffloadLedgerDeleteTest.java @@ -230,8 +230,16 @@ public void isOffloadedNeedsDeleteTest() throws Exception { .setBookkeeperDeleted(false) .build(); Boolean needsDelete = (Boolean) method.invoke(managedLedger, offloadContext); + Assert.assertFalse(needsDelete); + + offloadPolicies.setManagedLedgerOffloadDeletionLagInMillis(500L); + needsDelete = (Boolean) method.invoke(managedLedger, offloadContext); Assert.assertTrue(needsDelete); + offloadPolicies.setManagedLedgerOffloadDeletionLagInMillis(1000L * 2); + needsDelete = (Boolean) method.invoke(managedLedger, offloadContext); + Assert.assertFalse(needsDelete); + offloadContext = MLDataFormats.OffloadContext.newBuilder() .setTimestamp(config.getClock().millis() - 1000) .setComplete(false) @@ -248,13 +256,5 @@ public void isOffloadedNeedsDeleteTest() throws Exception { needsDelete = (Boolean) method.invoke(managedLedger, offloadContext); Assert.assertFalse(needsDelete); - offloadPolicies.setManagedLedgerOffloadDeletionLagInMillis(1000L * 2); - offloadContext = MLDataFormats.OffloadContext.newBuilder() - .setTimestamp(config.getClock().millis() - 1000) - .setComplete(true) - .setBookkeeperDeleted(false) - .build(); - needsDelete = (Boolean) method.invoke(managedLedger, offloadContext); - Assert.assertFalse(needsDelete); } }