diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java index 0455f0efa8bb6..82395b7002cb5 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java @@ -620,7 +620,8 @@ default void unfenceForInterceptorException() { void readyToCreateNewLedger(); /** - * Returns managed-ledger's properties. + * Returns a snapshot of the managed-ledger's properties. + * Changes made to the returned map are not applied to the managed ledger; use the property update methods instead. * * @return key-values of properties */ 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 29a68d32bdd20..f8686ce9b57f7 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 @@ -172,7 +172,7 @@ public Logger getLogger() { private final AtomicReference cacheEvictionPosition = new AtomicReference<>(); protected ManagedLedgerConfig config; - protected Map propertiesMap; + protected volatile Map propertiesMap; protected final MetaStore store; final ConcurrentLongHashMap> ledgerCache = @@ -445,17 +445,17 @@ public void operationComplete(ManagedLedgerInfo mlInfo, Stat stat) { ledgers.put(ls.getLedgerId(), ls); } - if (mlInfo.getPropertiesCount() > 0) { - propertiesMap = new HashMap<>(); - for (int i = 0; i < mlInfo.getPropertiesCount(); i++) { - KeyValue property = mlInfo.getPropertyAt(i); - propertiesMap.put(property.getKey(), property.getValue()); - } + Map loadedProperties = new ConcurrentHashMap<>(); + for (int i = 0; i < mlInfo.getPropertiesCount(); i++) { + KeyValue property = mlInfo.getPropertyAt(i); + loadedProperties.put(property.getKey(), property.getValue()); } - migrated = mlInfo.hasTerminatedPosition() && propertiesMap.containsKey(MIGRATION_STATE_PROPERTY); + migrated = mlInfo.hasTerminatedPosition() + && loadedProperties.containsKey(MIGRATION_STATE_PROPERTY); if (managedLedgerInterceptor != null) { - managedLedgerInterceptor.onManagedLedgerPropertiesInitialize(propertiesMap); + managedLedgerInterceptor.onManagedLedgerPropertiesInitialize(loadedProperties); } + propertiesMap = loadedProperties; // Last ledger stat may be zeroed, we must update it if (!ledgers.isEmpty()) { @@ -1398,8 +1398,22 @@ private long consumedLedgerSize(long ledgerSize, long ledgerEntries, long consum } public CompletableFuture asyncMigrate() { - propertiesMap.put(MIGRATION_STATE_PROPERTY, Boolean.TRUE.toString()); CompletableFuture result = new CompletableFuture<>(); + asyncSetProperty(MIGRATION_STATE_PROPERTY, Boolean.TRUE.toString(), new UpdatePropertiesCallback() { + @Override + public void updatePropertiesComplete(Map properties, Object ctx) { + terminateForMigration(result); + } + + @Override + public void updatePropertiesFailed(ManagedLedgerException exception, Object ctx) { + result.completeExceptionally(exception); + } + }, null); + return result; + } + + private void terminateForMigration(CompletableFuture result) { asyncTerminate(new TerminateCallback() { @Override @@ -1415,7 +1429,6 @@ public void terminateFailed(ManagedLedgerException exception, Object ctx) { result.completeExceptionally(exception); } }, null); - return result; } @Override @@ -4462,21 +4475,31 @@ private ManagedLedgerInfo getManagedLedgerInfo(LedgerInfo newLedger) { return buildManagedLedgerInfo(mlInfo); } private ManagedLedgerInfo buildManagedLedgerInfo(Map ledgers) { + return buildManagedLedgerInfo(ledgers, propertiesMap); + } + + private ManagedLedgerInfo buildManagedLedgerInfo(Map ledgers, + Map properties) { ManagedLedgerInfo mlInfo = new ManagedLedgerInfo(); mlInfo.addAllLedgerInfos(ledgers.values()); - return buildManagedLedgerInfo(mlInfo); + return buildManagedLedgerInfo(mlInfo, properties); } private ManagedLedgerInfo buildManagedLedgerInfo(ManagedLedgerInfo mlInfo) { + return buildManagedLedgerInfo(mlInfo, propertiesMap); + } + + private ManagedLedgerInfo buildManagedLedgerInfo(ManagedLedgerInfo mlInfo, + Map properties) { if (state == State.Terminated) { mlInfo.setTerminatedPosition() .setLedgerId(lastConfirmedEntry.getLedgerId()) .setEntryId(lastConfirmedEntry.getEntryId()); } if (managedLedgerInterceptor != null) { - managedLedgerInterceptor.onUpdateManagedLedgerInfo(propertiesMap); + managedLedgerInterceptor.onUpdateManagedLedgerInfo(properties); } - for (Map.Entry property : propertiesMap.entrySet()) { + for (Map.Entry property : properties.entrySet()) { mlInfo.addProperty().setKey(property.getKey()).setValue(property.getValue()); } @@ -4809,7 +4832,7 @@ public long getLastOffloadedFailureTimestamp() { @Override public Map getProperties() { - return propertiesMap; + return new HashMap<>(propertiesMap); } @Override @@ -4838,13 +4861,13 @@ public void asyncDeleteProperty(String key, final UpdatePropertiesCallback callb @Override public void setProperties(Map properties) throws InterruptedException, ManagedLedgerException { - updateProperties(properties, false, null); + updateProperties(new HashMap<>(properties), false, null); } @Override public void asyncSetProperties(Map properties, final UpdatePropertiesCallback callback, Object ctx) { - asyncUpdateProperties(properties, false, null, callback, ctx); + asyncUpdateProperties(new HashMap<>(properties), false, null, callback, ctx); } private void updateProperties(Map properties, boolean isDelete, @@ -4885,25 +4908,50 @@ private void asyncUpdateProperties(Map properties, boolean isDel callback, ctx), 100, TimeUnit.MILLISECONDS); return; } + Map updatedProperties = new HashMap<>(propertiesMap); if (isDelete) { - propertiesMap.remove(deleteKey); + updatedProperties.remove(deleteKey); } else { - propertiesMap.putAll(properties); + updatedProperties.putAll(properties); + } + + final ManagedLedgerInfo managedLedgerInfo; + final Map propertiesSnapshot; + try { + managedLedgerInfo = buildManagedLedgerInfo(ledgers, updatedProperties); + propertiesSnapshot = new ConcurrentHashMap<>(updatedProperties); + } catch (Throwable t) { + metadataMutex.unlock(); + callback.updatePropertiesFailed(ManagedLedgerException.getManagedLedgerException(t), ctx); + return; } - store.asyncUpdateLedgerIds(name, getManagedLedgerInfo(), ledgersStat, new MetaStoreCallback() { + + store.asyncUpdateLedgerIds(name, managedLedgerInfo, ledgersStat, new MetaStoreCallback() { @Override public void operationComplete(Void result, Stat version) { ledgersStat = version; - callback.updatePropertiesComplete(propertiesMap, ctx); + propertiesMap = propertiesSnapshot; metadataMutex.unlock(); + try { + callback.updatePropertiesComplete(new HashMap<>(propertiesSnapshot), ctx); + } catch (Throwable t) { + log.error().exception(t).log("Managed ledger properties callback failed"); + } } @Override public void operationFailed(MetaStoreException e) { log.error().exception(e).log("Update managedLedger's properties failed"); - handleBadVersion(e); - callback.updatePropertiesFailed(e, ctx); - metadataMutex.unlock(); + try { + handleBadVersion(e); + } finally { + metadataMutex.unlock(); + } + try { + callback.updatePropertiesFailed(e, ctx); + } catch (Throwable t) { + log.error().exception(t).log("Managed ledger properties callback failed"); + } } }); } diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java index ab0a0411dfc7f..ad806ae121cd8 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java @@ -1690,6 +1690,257 @@ public void updatePropertiesFailed(ManagedLedgerException exception, Object ctx) assertTrue(latch3.await(5, TimeUnit.SECONDS)); } + @Test + public void testPropertiesSnapshotsAreDetached() throws Exception { + ManagedLedger ledger = factory.open("properties-snapshot-test"); + ledger.setProperties(Map.of("key1", "value1")); + + Map firstSnapshot = ledger.getProperties(); + firstSnapshot.put("external", "value"); + assertEquals(ledger.getProperties(), Map.of("key1", "value1")); + + ledger.close(); + ledger = factory.open("properties-snapshot-test"); + assertEquals(ledger.getProperties(), Map.of("key1", "value1")); + + CountDownLatch callbackCompleted = new CountDownLatch(1); + ledger.asyncSetProperty("key2", "value2", new AsyncCallbacks.UpdatePropertiesCallback() { + @Override + public void updatePropertiesComplete(Map properties, Object ctx) { + properties.put("callback", "value"); + callbackCompleted.countDown(); + } + + @Override + public void updatePropertiesFailed(ManagedLedgerException exception, Object ctx) { + callbackCompleted.countDown(); + } + }, null); + + assertTrue(callbackCompleted.await(5, TimeUnit.SECONDS)); + assertEquals(ledger.getProperties(), Map.of("key1", "value1", "key2", "value2")); + assertEquals(firstSnapshot, Map.of("key1", "value1", "external", "value")); + } + + @Test + public void testFailedPropertiesUpdateDoesNotChangeInMemoryState() throws Exception { + String ledgerName = "properties-failure-test"; + ManagedLedger ledger = factory.open(ledgerName); + ledger.setProperty("existing", "value"); + + metadataStore.failConditional(new MetadataStoreException("injected failure"), + (operation, path) -> operation == FaultInjectionMetadataStore.OperationType.PUT + && path.equals("/managed-ledgers/" + ledgerName)); + + CompletableFuture failure = new CompletableFuture<>(); + ledger.asyncSetProperty("failed", "value", new AsyncCallbacks.UpdatePropertiesCallback() { + @Override + public void updatePropertiesComplete(Map properties, Object ctx) { + failure.completeExceptionally(new AssertionError("The properties update should fail")); + } + + @Override + public void updatePropertiesFailed(ManagedLedgerException exception, Object ctx) { + failure.complete(exception); + } + }, null); + + assertNotNull(failure.get(5, TimeUnit.SECONDS)); + assertEquals(ledger.getProperties(), Map.of("existing", "value")); + + ledger.setProperty("after-failure", "value"); + assertEquals(ledger.getProperties(), Map.of("existing", "value", "after-failure", "value")); + } + + @Test + @SuppressWarnings("unchecked") + public void testPropertiesRemainUnchangedUntilMetadataUpdateCompletes() throws Exception { + String ledgerName = "properties-pending-update-test"; + String ledgerPath = "/managed-ledgers/" + ledgerName; + CompletableFuture releaseUpdate = new CompletableFuture<>(); + CountDownLatch updateIntercepted = new CountDownLatch(1); + AtomicBoolean interceptNextPut = new AtomicBoolean(false); + + FaultInjectionMetadataStore spyStore = spy(metadataStore); + doAnswer(invocation -> { + if (ledgerPath.equals(invocation.getArgument(0)) + && interceptNextPut.compareAndSet(true, false)) { + updateIntercepted.countDown(); + CompletableFuture realResult = (CompletableFuture) invocation.callRealMethod(); + CompletableFuture gatedResult = new CompletableFuture<>(); + realResult.whenComplete((stat, error) -> releaseUpdate.whenComplete((ignored, releaseError) -> { + if (error != null) { + gatedResult.completeExceptionally(error); + } else { + gatedResult.complete(stat); + } + })); + return gatedResult; + } + return invocation.callRealMethod(); + }).when(spyStore).put(eq(ledgerPath), any(byte[].class), any()); + + ManagedLedgerFactoryImpl gatedFactory = new ManagedLedgerFactoryImpl(spyStore, bkc); + try { + ManagedLedger ledger = gatedFactory.open(ledgerName); + ledger.setProperties(Map.of("key1", "old", "key2", "old")); + + Map update = new HashMap<>(); + update.put("key1", "new"); + update.put("key2", "new"); + CompletableFuture updateResult = new CompletableFuture<>(); + interceptNextPut.set(true); + ledger.asyncSetProperties(update, new AsyncCallbacks.UpdatePropertiesCallback() { + @Override + public void updatePropertiesComplete(Map properties, Object ctx) { + updateResult.complete(null); + } + + @Override + public void updatePropertiesFailed(ManagedLedgerException exception, Object ctx) { + updateResult.completeExceptionally(exception); + } + }, null); + + assertTrue(updateIntercepted.await(5, TimeUnit.SECONDS)); + update.put("late-change", "value"); + assertEquals(ledger.getProperties(), Map.of("key1", "old", "key2", "old")); + + releaseUpdate.complete(null); + updateResult.get(5, TimeUnit.SECONDS); + assertEquals(ledger.getProperties(), Map.of("key1", "new", "key2", "new")); + } finally { + releaseUpdate.complete(null); + gatedFactory.shutdown(); + } + } + + @Test + public void testPropertiesCallbackExceptionDoesNotKeepMetadataMutexLocked() throws Exception { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("properties-callback-exception-test"); + CountDownLatch firstCallbackInvoked = new CountDownLatch(1); + ledger.asyncSetProperty("key1", "value1", new AsyncCallbacks.UpdatePropertiesCallback() { + @Override + public void updatePropertiesComplete(Map properties, Object ctx) { + firstCallbackInvoked.countDown(); + throw new RuntimeException("injected callback failure"); + } + + @Override + public void updatePropertiesFailed(ManagedLedgerException exception, Object ctx) { + firstCallbackInvoked.countDown(); + } + }, null); + + assertTrue(firstCallbackInvoked.await(5, TimeUnit.SECONDS)); + + CountDownLatch secondCallbackCompleted = new CountDownLatch(1); + AtomicReference secondFailure = new AtomicReference<>(); + ledger.asyncSetProperty("key2", "value2", new AsyncCallbacks.UpdatePropertiesCallback() { + @Override + public void updatePropertiesComplete(Map properties, Object ctx) { + secondCallbackCompleted.countDown(); + } + + @Override + public void updatePropertiesFailed(ManagedLedgerException exception, Object ctx) { + secondFailure.set(exception); + secondCallbackCompleted.countDown(); + } + }, null); + + assertTrue(secondCallbackCompleted.await(5, TimeUnit.SECONDS)); + assertNull(secondFailure.get()); + assertEquals(ledger.getProperties(), Map.of("key1", "value1", "key2", "value2")); + assertTrue(ledger.metadataMutex.tryLock()); + assertFalse(ledger.metadataMutex.tryLock()); + ledger.metadataMutex.unlock(); + } + + @Test + @SuppressWarnings("unchecked") + public void testConcurrentPropertyUpdateDoesNotEraseMigrationMarker() throws Exception { + String ledgerName = "properties-migration-race-test"; + String ledgerPath = "/managed-ledgers/" + ledgerName; + CompletableFuture releasePropertyUpdate = new CompletableFuture<>(); + CountDownLatch propertyUpdateIntercepted = new CountDownLatch(1); + AtomicBoolean interceptNextPut = new AtomicBoolean(false); + + FaultInjectionMetadataStore spyStore = spy(metadataStore); + doAnswer(invocation -> { + if (ledgerPath.equals(invocation.getArgument(0)) + && interceptNextPut.compareAndSet(true, false)) { + propertyUpdateIntercepted.countDown(); + CompletableFuture realResult = (CompletableFuture) invocation.callRealMethod(); + CompletableFuture gatedResult = new CompletableFuture<>(); + realResult.whenComplete((stat, error) -> releasePropertyUpdate.whenComplete((ignored, releaseError) -> { + if (error != null) { + gatedResult.completeExceptionally(error); + } else { + gatedResult.complete(stat); + } + })); + return gatedResult; + } + return invocation.callRealMethod(); + }).when(spyStore).put(eq(ledgerPath), any(byte[].class), any()); + + ManagedLedgerFactoryImpl gatedFactory = new ManagedLedgerFactoryImpl(spyStore, bkc); + CompletableFuture releaseLedgerClose = new CompletableFuture<>(); + try { + ManagedLedgerImpl ledger = (ManagedLedgerImpl) gatedFactory.open(ledgerName); + ledger.setProperty("existing", "value"); + + LedgerHandle currentLedger = ledger.currentLedger; + LedgerHandle spyLedgerHandle = spy(currentLedger); + CountDownLatch ledgerCloseIntercepted = new CountDownLatch(1); + doAnswer(invocation -> { + AsyncCallback.CloseCallback callback = invocation.getArgument(0); + Object closeContext = invocation.getArgument(1); + currentLedger.asyncClose((rc, lh, ctx) -> { + ledgerCloseIntercepted.countDown(); + releaseLedgerClose.whenComplete((ignored, error) -> + callback.closeComplete(rc, spyLedgerHandle, closeContext)); + }, closeContext); + return null; + }).when(spyLedgerHandle).asyncClose(any(AsyncCallback.CloseCallback.class), any()); + ledger.currentLedger = spyLedgerHandle; + + CompletableFuture propertyResult = new CompletableFuture<>(); + interceptNextPut.set(true); + ledger.asyncSetProperty("pending", "value", new AsyncCallbacks.UpdatePropertiesCallback() { + @Override + public void updatePropertiesComplete(Map properties, Object ctx) { + propertyResult.complete(null); + } + + @Override + public void updatePropertiesFailed(ManagedLedgerException exception, Object ctx) { + propertyResult.completeExceptionally(exception); + } + }, null); + + assertTrue(propertyUpdateIntercepted.await(5, TimeUnit.SECONDS)); + CompletableFuture migrationResult = ledger.asyncMigrate(); + + releasePropertyUpdate.complete(null); + propertyResult.get(5, TimeUnit.SECONDS); + assertTrue(ledgerCloseIntercepted.await(5, TimeUnit.SECONDS)); + releaseLedgerClose.complete(null); + migrationResult.get(5, TimeUnit.SECONDS); + + ledger.close(); + ManagedLedger reopenedLedger = gatedFactory.open(ledgerName); + assertTrue(reopenedLedger.isMigrated()); + assertEquals(reopenedLedger.getProperties(), + Map.of("existing", "value", "pending", "value", "migrated", "true")); + } finally { + releasePropertyUpdate.complete(null); + releaseLedgerClose.complete(null); + gatedFactory.shutdown(); + } + } + @Test public void testConcurrentAsyncSetProperties() throws Exception { final CountDownLatch latch = new CountDownLatch(1000);