From 5f60f6d2de5a17632751d19479e7a3d1b375571c Mon Sep 17 00:00:00 2001 From: coderzc Date: Mon, 19 Sep 2022 12:35:11 +0800 Subject: [PATCH 1/3] Support internal cursor properties. --- .../bookkeeper/mledger/ManagedCursor.java | 4 ++- .../mledger/impl/ManagedCursorImpl.java | 25 ++++++++++++++++- .../impl/ManagedCursorPropertiesTest.java | 28 +++++++++++++++++++ 3 files changed, 55 insertions(+), 2 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java index 17dbac09a2292..41833bf75f444 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java @@ -97,7 +97,9 @@ enum IndividualDeletedEntries { CompletableFuture putCursorProperty(String key, String value); /** - * Set all properties associated with the cursor. + * Set all properties associated with the cursor, + * but internal properties(start with '#pulsar_internal.') are still preserved + * and prohibit setting of internal properties by this method. * * Note: {@link ManagedLedgerException.BadVersionException} will be set in this {@link CompletableFuture}, * if there are concurrent modification and store data has changed. diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 6d595e76dc127..81f0d28c8d8b6 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -121,6 +121,8 @@ public class ManagedCursorImpl implements ManagedCursor { protected final ManagedLedgerImpl ledger; private final String name; + public static final String CURSOR_INTERNAL_PROPERTY_PREFIX = "#pulsar.internal."; + private volatile Map cursorProperties; private final BookKeeper.DigestType digestType; @@ -371,7 +373,28 @@ public void operationFailed(MetaStoreException e) { @Override public CompletableFuture setCursorProperties(Map cursorProperties) { - return computeCursorProperties(lastRead -> cursorProperties); + Map newProperties = + cursorProperties == null ? new HashMap<>() : new HashMap<>(cursorProperties); + + // Prohibit setting of internal properties + Set keys = newProperties.keySet(); + for (String key : keys) { + if (key.startsWith(CURSOR_INTERNAL_PROPERTY_PREFIX)) { + throw new IllegalArgumentException( + "The property key can't start with " + CURSOR_INTERNAL_PROPERTY_PREFIX); + } + } + + return computeCursorProperties(lastRead -> { + if (lastRead != null) { + lastRead.forEach((k, v) -> { + if (k.startsWith(CURSOR_INTERNAL_PROPERTY_PREFIX)) { + newProperties.put(k, v); + } + }); + } + return newProperties; + }); } @Override diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java index c9082344d23a8..0585423417021 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java @@ -18,9 +18,11 @@ */ package org.apache.bookkeeper.mledger.impl; +import static org.apache.bookkeeper.mledger.impl.ManagedCursorImpl.CURSOR_INTERNAL_PROPERTY_PREFIX; import static org.apache.bookkeeper.mledger.util.Futures.executeWithRetry; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertNull; +import static org.testng.Assert.assertTrue; import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; @@ -37,6 +39,7 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.test.MockedBookKeeperTestCase; import org.apache.pulsar.common.api.proto.CommandSubscribe.InitialPosition; +import org.testng.Assert; import org.testng.annotations.Test; public class ManagedCursorPropertiesTest extends MockedBookKeeperTestCase { @@ -232,6 +235,31 @@ void testUpdateCursorProperties() throws Exception { ledger.close(); factory2.shutdown(); + + // Create a new factory to force a managed ledger close and recovery + ManagedLedgerFactory factory3 = new ManagedLedgerFactoryImpl(metadataStore, bkc); + // Reopen the managed ledger + ledger = factory3.open("testUpdateCursorProperties", new ManagedLedgerConfig()); + c1 = ledger.openCursor("c1"); + + c1.putCursorProperty(CURSOR_INTERNAL_PROPERTY_PREFIX + "test", "test").get(); + c1.putCursorProperty("custom4", "custom4").get(); + c1.setCursorProperties(cursorPropertiesUpdated).get(); + + cursorPropertiesUpdated.put(CURSOR_INTERNAL_PROPERTY_PREFIX + "test", "test"); + + try { + c1.setCursorProperties(cursorPropertiesUpdated).get(); + Assert.fail("Should fail"); + } catch (IllegalArgumentException e) { + assertTrue(e.getMessage().contains("The property key can't start with")); + } + + assertEquals(c1.getCursorProperties(), cursorPropertiesUpdated); + + ledger.close(); + + factory3.shutdown(); } @Test From d937eaa2516eae723ea754d8e8f76959f997f2e4 Mon Sep 17 00:00:00 2001 From: coderzc Date: Mon, 7 Nov 2022 16:19:05 +0800 Subject: [PATCH 2/3] Address comment --- .../bookkeeper/mledger/impl/ManagedCursorImpl.java | 5 +++-- .../mledger/impl/ManagedCursorPropertiesTest.java | 10 +++++----- 2 files changed, 8 insertions(+), 7 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 81f0d28c8d8b6..63cda621dbedb 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -96,6 +96,7 @@ import org.apache.bookkeeper.mledger.proto.MLDataFormats.PositionInfo; import org.apache.bookkeeper.mledger.proto.MLDataFormats.StringProperty; import org.apache.commons.lang3.tuple.Pair; +import org.apache.pulsar.common.util.FutureUtil; import org.apache.pulsar.common.util.collections.BitSetRecyclable; import org.apache.pulsar.common.util.collections.LongPairRangeSet; import org.apache.pulsar.common.util.collections.LongPairRangeSet.LongPairConsumer; @@ -380,8 +381,8 @@ public CompletableFuture setCursorProperties(Map cursorPro Set keys = newProperties.keySet(); for (String key : keys) { if (key.startsWith(CURSOR_INTERNAL_PROPERTY_PREFIX)) { - throw new IllegalArgumentException( - "The property key can't start with " + CURSOR_INTERNAL_PROPERTY_PREFIX); + return FutureUtil.failedFuture(new IllegalArgumentException( + "The property key can't start with " + CURSOR_INTERNAL_PROPERTY_PREFIX)); } } diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java index 0585423417021..16f6c462f8572 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java @@ -242,14 +242,14 @@ void testUpdateCursorProperties() throws Exception { ledger = factory3.open("testUpdateCursorProperties", new ManagedLedgerConfig()); c1 = ledger.openCursor("c1"); - c1.putCursorProperty(CURSOR_INTERNAL_PROPERTY_PREFIX + "test", "test").get(); - c1.putCursorProperty("custom4", "custom4").get(); - c1.setCursorProperties(cursorPropertiesUpdated).get(); + c1.putCursorProperty(CURSOR_INTERNAL_PROPERTY_PREFIX + "test", "test").get(10, TimeUnit.SECONDS); + c1.putCursorProperty("custom4", "custom4").get(10, TimeUnit.SECONDS); + c1.setCursorProperties(cursorPropertiesUpdated).get(10, TimeUnit.SECONDS); cursorPropertiesUpdated.put(CURSOR_INTERNAL_PROPERTY_PREFIX + "test", "test"); try { - c1.setCursorProperties(cursorPropertiesUpdated).get(); + c1.setCursorProperties(cursorPropertiesUpdated).get(10, TimeUnit.SECONDS); Assert.fail("Should fail"); } catch (IllegalArgumentException e) { assertTrue(e.getMessage().contains("The property key can't start with")); @@ -284,7 +284,7 @@ public void testUpdateCursorPropertiesConcurrent() throws Exception { ManagedLedgerException.BadVersionException.class)); for (CompletableFuture future : futures) { - future.get(); + future.get(10, TimeUnit.SECONDS); } assertEquals(c1.getCursorProperties().get("a"), "2"); From 3a5056acb983605416e3d644d80e6f49ce2fcf8d Mon Sep 17 00:00:00 2001 From: coderzc Date: Mon, 7 Nov 2022 17:20:15 +0800 Subject: [PATCH 3/3] fix test --- .../mledger/impl/ManagedCursorPropertiesTest.java | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java index 16f6c462f8572..23c7142eb9b9b 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java @@ -39,6 +39,7 @@ import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.test.MockedBookKeeperTestCase; import org.apache.pulsar.common.api.proto.CommandSubscribe.InitialPosition; +import org.apache.pulsar.common.util.FutureUtil; import org.testng.Assert; import org.testng.annotations.Test; @@ -251,8 +252,9 @@ void testUpdateCursorProperties() throws Exception { try { c1.setCursorProperties(cursorPropertiesUpdated).get(10, TimeUnit.SECONDS); Assert.fail("Should fail"); - } catch (IllegalArgumentException e) { - assertTrue(e.getMessage().contains("The property key can't start with")); + } catch (Exception e) { + assertTrue( + FutureUtil.unwrapCompletionException(e).getMessage().contains("The property key can't start with")); } assertEquals(c1.getCursorProperties(), cursorPropertiesUpdated);