-
Notifications
You must be signed in to change notification settings - Fork 3
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
20 changed files
with
624 additions
and
69 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
32 changes: 25 additions & 7 deletions
32
src/main/java/uk/gov/justice/artemis/manager/connector/combined/CombinedManagement.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,38 +1,56 @@ | ||
package uk.gov.justice.artemis.manager.connector.combined; | ||
|
||
import uk.gov.justice.artemis.manager.connector.combined.duplicate.AddedMessageFinder; | ||
import uk.gov.justice.artemis.manager.connector.combined.duplicate.BrowsedMessages; | ||
import uk.gov.justice.artemis.manager.connector.combined.duplicate.DuplicateMessageFinder; | ||
import uk.gov.justice.artemis.manager.connector.combined.duplicate.DuplicateMessageRemover; | ||
import uk.gov.justice.artemis.manager.connector.combined.duplicate.DuplicateMessages; | ||
import uk.gov.justice.artemis.manager.connector.combined.duplicate.TopicDuplicateMessageFinder; | ||
|
||
import java.util.List; | ||
|
||
public class CombinedManagement { | ||
|
||
private final DuplicateMessageFinder duplicateMessageFinder; | ||
private final TopicDuplicateMessageFinder topicDuplicateMessageFinder; | ||
private final DuplicateMessageRemover duplicateMessageRemover; | ||
private final AddedMessageFinder addedMessageFinder; | ||
|
||
public CombinedManagement(final DuplicateMessageFinder duplicateMessageFinder, | ||
final TopicDuplicateMessageFinder topicDuplicateMessageFinder, | ||
final DuplicateMessageRemover duplicateMessageRemover, | ||
final AddedMessageFinder addedMessageFinder) { | ||
this.duplicateMessageFinder = duplicateMessageFinder; | ||
this.topicDuplicateMessageFinder = topicDuplicateMessageFinder; | ||
this.duplicateMessageRemover = duplicateMessageRemover; | ||
this.addedMessageFinder = addedMessageFinder; | ||
} | ||
|
||
public CombinedManagementFunction<List<String>> removeAllDuplicates() { | ||
|
||
return (queueBrowser, queueSender, jmsQueueControl) -> | ||
{ | ||
final DuplicateMessages duplicateMessages = duplicateMessageFinder.findDuplicateMessages(queueBrowser); | ||
return (queueBrowser, queueSender, jmsQueueControl) -> { | ||
|
||
duplicateMessageRemover.removeDuplicatesOnly( | ||
final BrowsedMessages browsedMessages = duplicateMessageFinder.findDuplicateMessages(queueBrowser); | ||
|
||
duplicateMessageRemover.removeAndResendDuplicateMessages( | ||
queueSender, | ||
jmsQueueControl, | ||
browsedMessages); | ||
|
||
return addedMessageFinder.findAddedMessages(browsedMessages, queueBrowser); | ||
}; | ||
} | ||
|
||
public CombinedManagementFunction<List<String>> deduplicateTopicMessages() { | ||
return (queueBrowser, queueSender, jmsQueueControl) -> { | ||
|
||
final BrowsedMessages topicBrowsedMessages = topicDuplicateMessageFinder.findTopicDuplicateMessages(queueBrowser); | ||
|
||
duplicateMessageRemover.removeAndResendDuplicateMessages( | ||
queueSender, | ||
jmsQueueControl, | ||
duplicateMessages); | ||
topicBrowsedMessages); | ||
|
||
return addedMessageFinder.findAddedMessages(duplicateMessages, queueBrowser); | ||
return addedMessageFinder.findAddedMessages(topicBrowsedMessages, queueBrowser); | ||
}; | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
11 changes: 6 additions & 5 deletions
11
...combined/duplicate/DuplicateMessages.java → ...r/combined/duplicate/BrowsedMessages.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
70 changes: 70 additions & 0 deletions
70
...gov/justice/artemis/manager/connector/combined/duplicate/TopicDuplicateMessageFinder.java
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,70 @@ | ||
package uk.gov.justice.artemis.manager.connector.combined.duplicate; | ||
|
||
import uk.gov.justice.artemis.manager.connector.combined.CombinedManagementFunctionException; | ||
|
||
import java.util.ArrayList; | ||
import java.util.Enumeration; | ||
import java.util.HashMap; | ||
import java.util.List; | ||
import java.util.Map; | ||
|
||
import javax.jms.JMSException; | ||
import javax.jms.Message; | ||
import javax.jms.QueueBrowser; | ||
|
||
public class TopicDuplicateMessageFinder { | ||
|
||
private final JmsMessageUtil jmsMessageUtil; | ||
|
||
public TopicDuplicateMessageFinder(final JmsMessageUtil jmsMessageUtil) { | ||
this.jmsMessageUtil = jmsMessageUtil; | ||
} | ||
|
||
@SuppressWarnings("unchecked") | ||
public BrowsedMessages findTopicDuplicateMessages(final QueueBrowser queueBrowser) { | ||
|
||
final Map<String, List<Message>> duplicateMessages = new HashMap<>(); | ||
final Map<String, Message> messageCache = new HashMap<>(); | ||
|
||
try { | ||
final Enumeration<Message> browserEnumeration = queueBrowser.getEnumeration(); | ||
|
||
while (browserEnumeration.hasMoreElements()) { | ||
|
||
final Message message = browserEnumeration.nextElement(); | ||
final String jmsMessageID = jmsMessageUtil.getJmsMessageIdFrom(message); | ||
|
||
if (messageCache.containsKey(jmsMessageID)) { | ||
|
||
final Message originalMessage = messageCache.get(jmsMessageID); | ||
final String originalConsumer = jmsMessageUtil.getConsumerFrom(originalMessage); | ||
final String duplicateConsumer = jmsMessageUtil.getConsumerFrom(message); | ||
final boolean isDuplicateTopicMessage = !duplicateConsumer.equals(originalConsumer); | ||
|
||
if (isDuplicateTopicMessage) { | ||
|
||
duplicateMessages | ||
.computeIfAbsent(jmsMessageID, key -> createNewMessageListWith(originalMessage)) | ||
.add(message); | ||
} else { | ||
messageCache.put(jmsMessageID, message); | ||
} | ||
|
||
} else { | ||
messageCache.put(jmsMessageID, message); | ||
} | ||
} | ||
|
||
return new BrowsedMessages(duplicateMessages, messageCache); | ||
|
||
} catch (final JMSException exception) { | ||
throw new CombinedManagementFunctionException("Failed to browse messages on queue.", exception); | ||
} | ||
} | ||
|
||
private List<Message> createNewMessageListWith(final Message originalMessage) { | ||
final ArrayList<Message> messages = new ArrayList<>(); | ||
messages.add(originalMessage); | ||
return messages; | ||
} | ||
} |
Oops, something went wrong.