Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -328,6 +328,14 @@ protected Policies getNamespacePolicies(NamespaceName namespaceName) {

}

protected PersistencePolicies getOrDefaultPersistencePolicy(PersistencePolicies policies) {
return policies != null ? policies
: new PersistencePolicies(pulsar().getConfiguration().getManagedLedgerDefaultEnsembleSize(),
pulsar().getConfiguration().getManagedLedgerDefaultWriteQuorum(),
pulsar().getConfiguration().getManagedLedgerDefaultAckQuorum(),
pulsar().getConfiguration().getManagedLedgerDefaultMarkDeleteRateLimit());
}

protected CompletableFuture<Policies> getNamespacePoliciesAsync(NamespaceName namespaceName) {
return namespaceResources().getPoliciesAsync(namespaceName).thenCompose(policies -> {
if (policies.isPresent()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1273,6 +1273,14 @@ protected CompletableFuture<Void> internalSetPersistenceAsync(PersistencePolicie
.thenCompose(__ -> doUpdatePersistenceAsync(persistence));
}


protected PersistencePolicies internalGetPersistence() {
validateNamespacePolicyOperation(namespaceName, PolicyName.PERSISTENCE, PolicyOperation.READ);

Policies policies = getNamespacePolicies(namespaceName);
return getOrDefaultPersistencePolicy(policies.persistence);
}

private CompletableFuture<Void> doUpdatePersistenceAsync(PersistencePolicies persistence) {
return updatePoliciesAsync(namespaceName, policies -> {
policies.persistence = persistence;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3517,13 +3517,7 @@ protected CompletableFuture<PersistencePolicies> internalGetPersistence(boolean
if (applied) {
PersistencePolicies namespacePolicy = getNamespacePolicies(namespaceName)
.persistence;
return namespacePolicy == null
? new PersistencePolicies(
pulsar().getConfiguration().getManagedLedgerDefaultEnsembleSize(),
pulsar().getConfiguration().getManagedLedgerDefaultWriteQuorum(),
pulsar().getConfiguration().getManagedLedgerDefaultAckQuorum(),
pulsar().getConfiguration().getManagedLedgerDefaultMarkDeleteRateLimit())
: namespacePolicy;
return getOrDefaultPersistencePolicy(namespacePolicy);
}
return null;
}));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -436,7 +436,7 @@ public void testSetPersistencePolicies() throws Exception {
final String namespace = "prop-xyz/ns2";
admin.namespaces().createNamespace(namespace, Set.of("test"));

assertNull(admin.namespaces().getPersistence(namespace));
assertNotNull(admin.namespaces().getPersistence(namespace));
admin.namespaces().setPersistence(namespace, new PersistencePolicies(3, 3, 3, 10.0));
assertEquals(admin.namespaces().getPersistence(namespace), new PersistencePolicies(3, 3, 3, 10.0));

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import static org.mockito.Mockito.verify;
import static org.testng.Assert.assertEquals;
import static org.testng.Assert.assertFalse;
import static org.testng.Assert.assertEquals;
import static org.testng.Assert.assertNotEquals;
import static org.testng.Assert.assertNotNull;
import static org.testng.Assert.assertNull;
Expand Down Expand Up @@ -829,7 +830,7 @@ public void namespaces() throws Exception {
policies.is_allow_auto_update_schema = conf.isAllowAutoUpdateSchemaEnabled();
assertEquals(admin.namespaces().getPolicies("prop-xyz/ns1"), policies);

assertEquals(admin.namespaces().getPersistence("prop-xyz/ns1"), null);
assertEquals(admin.namespaces().getPersistence("prop-xyz/ns1"), new PersistencePolicies(2, 2, 2, 1.0));
admin.namespaces().setPersistence("prop-xyz/ns1", new PersistencePolicies(3, 2, 1, 10.0));
assertEquals(admin.namespaces().getPersistence("prop-xyz/ns1"), new PersistencePolicies(3, 2, 1, 10.0));

Expand Down Expand Up @@ -3453,6 +3454,17 @@ public void testPeekEncryptedMessages() throws Exception {
}
}

@Test
public void testGetPersistenceAPI() throws Exception {
String namespace = "prop-xyz/test/persistence";
admin.namespaces().createNamespace(namespace);
PersistencePolicies persistencePolicies = admin.namespaces().getPersistence(namespace);
assertNotNull(persistencePolicies);
assertEquals(persistencePolicies.getBookkeeperEnsemble(), conf.getManagedLedgerDefaultEnsembleSize());
assertEquals(persistencePolicies.getBookkeeperWriteQuorum(), conf.getManagedLedgerDefaultWriteQuorum());
assertEquals(persistencePolicies.getBookkeeperAckQuorum(), conf.getManagedLedgerDefaultAckQuorum());
}

@Test
public void testGetPartitionStatsWithEarliestTimeInBacklog() throws PulsarAdminException, PulsarClientException {
final String topicName = "persistent://prop-xyz/ns1/testPeekEncryptedMessages-" + UUID.randomUUID();
Expand Down