michaelandrepearce commented on a change in pull request #2850: ARTEMIS-2504
implement retroactive addresses
URL: https://github.com/apache/activemq-artemis/pull/2850#discussion_r328400809
##########
File path:
artemis-server/src/main/java/org/apache/activemq/artemis/core/postoffice/impl/PostOfficeImpl.java
##########
@@ -467,6 +478,83 @@ private boolean internalAddressInfo(AddressInfo
addressInfo, boolean reload) thr
}
}
+ private void registerRepositoryListenerForRetroactiveAddress(SimpleString
name) {
+ HierarchicalRepositoryChangeListener repositoryChangeListener = () -> {
+ String prefix = server.getConfiguration().getInternalNamingPrefix();
+ String address =
ResourceNames.decomposeRetroactiveResourceName(prefix, name.toString(),
ResourceNames.ADDRESS);
+ AddressSettings settings =
addressSettingsRepository.getMatch(address);
+ Queue internalQueue =
server.locateQueue(ResourceNames.getRetroactiveResourceName(prefix,
SimpleString.toSimpleString(address), ResourceNames.QUEUE));
+ if (internalQueue != null && internalQueue.getRingSize() !=
settings.getRetroactiveMessageCount()) {
+ internalQueue.setRingSize(settings.getRetroactiveMessageCount());
+ }
+ };
+ addressSettingsRepository.registerListener(repositoryChangeListener);
+
server.getAddressInfo(name).setRepositoryChangeListener(repositoryChangeListener);
+ }
+
+ private void createRetroactiveResources(final SimpleString
retroactiveAddressName, final long retroactiveMessageCount, final boolean
reload) throws Exception {
+ String prefix = server.getConfiguration().getInternalNamingPrefix();
+ final SimpleString internalAddressName =
ResourceNames.getRetroactiveResourceName(prefix, retroactiveAddressName,
ResourceNames.ADDRESS);
+ final SimpleString internalQueueName =
ResourceNames.getRetroactiveResourceName(prefix, retroactiveAddressName,
ResourceNames.QUEUE);
+ final SimpleString internalDivertName =
ResourceNames.getRetroactiveResourceName(prefix, retroactiveAddressName,
ResourceNames.DIVERT);
+
+ if (!reload) {
+ AddressInfo addressInfo = new
AddressInfo(internalAddressName).addRoutingType(RoutingType.ANYCAST).setInternal(true);
+ addAddressInfo(addressInfo);
+ server.createQueue(internalAddressName,
+ RoutingType.ANYCAST,
+ internalQueueName,
+ null,
+ null,
+ true,
+ false,
+ false,
+ false,
+ false,
+ 0,
+ false,
+ false,
+ false,
+ 0,
+ null,
+ false,
+ null,
+ false,
+ 0,
+ 0L,
+ false,
+ 0L,
+ 0L,
+ false,
+ retroactiveMessageCount);
+ }
+ server.deployDivert(new DivertConfiguration()
+ .setName(internalDivertName.toString())
+ .setAddress(retroactiveAddressName.toString())
+ .setExclusive(false)
+
.setForwardingAddress(internalAddressName.toString())
+
.setRoutingType(ComponentConfigurationRoutingType.STRIP));
Review comment:
It may mean in a retro anycast and a retro multicast queue is required to
preserve that
----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
For queries about this service, please contact Infrastructure at:
[email protected]
With regards,
Apache Git Services