Skip to content
Open
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 @@ -24,6 +24,7 @@
import java.util.Iterator;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.CopyOnWriteArrayList;
Expand Down Expand Up @@ -119,7 +120,7 @@ public class ActiveMQConnection implements Connection, TopicConnection, QueueCon

private static final Logger LOG = LoggerFactory.getLogger(ActiveMQConnection.class);

public final ConcurrentMap<ActiveMQTempDestination, ActiveMQTempDestination> activeTempDestinations = new ConcurrentHashMap<>();
final Set<ActiveMQTempDestination> activeTempDestinations = ConcurrentHashMap.newKeySet();

protected boolean dispatchAsync=true;
protected boolean alwaysSessionAsync = true;
Expand Down Expand Up @@ -2151,7 +2152,7 @@ protected ActiveMQTempDestination createTempDestination(boolean topic) throws JM
syncSendPacket(info);

dest.setConnection(this);
activeTempDestinations.put(dest, dest);
activeTempDestinations.add(dest);
return dest;
}

Expand Down Expand Up @@ -2187,7 +2188,7 @@ public boolean isDeleted(ActiveMQDestination dest) {
return false;
}

return !activeTempDestinations.containsValue(dest);
return !activeTempDestinations.contains(dest);
}

public boolean isCopyMessageOnSend() {
Expand Down Expand Up @@ -2575,21 +2576,17 @@ public void setRmIdFromConnectionId(boolean rmIdFromConnectionId) {
*/
public void cleanUpTempDestinations() {

if (this.activeTempDestinations == null || this.activeTempDestinations.isEmpty()) {
if (this.activeTempDestinations.isEmpty()) {
return;
}

Iterator<ConcurrentMap.Entry<ActiveMQTempDestination, ActiveMQTempDestination>> entries
= this.activeTempDestinations.entrySet().iterator();
while(entries.hasNext()) {
ConcurrentMap.Entry<ActiveMQTempDestination, ActiveMQTempDestination> entry = entries.next();
for (ActiveMQTempDestination dest : activeTempDestinations) {
try {
// Only delete this temp destination if it was created from this connection. The connection used
// for the advisory consumer may also have a reference to this temp destination.
ActiveMQTempDestination dest = entry.getValue();
String thisConnectionId = (info.getConnectionId() == null) ? "" : info.getConnectionId().toString();
if (dest.getConnectionId() != null && dest.getConnectionId().equals(thisConnectionId)) {
this.deleteTempDestination(entry.getValue());
this.deleteTempDestination(dest);
}
} catch (Exception ex) {
// the temp dest is in use so it can not be deleted.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -100,7 +100,7 @@ private void processDestinationInfo(DestinationInfo dinfo) {
if (tempDest.getConnection() != null) {
tempDest = (ActiveMQTempDestination) tempDest.createDestination(tempDest.getPhysicalName());
}
connection.activeTempDestinations.put(tempDest, tempDest);
connection.activeTempDestinations.add(tempDest);
} else if (dinfo.getOperationType() == DestinationInfo.REMOVE_OPERATION_TYPE) {
connection.activeTempDestinations.remove(tempDest);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -221,12 +221,8 @@ public void testPublishFailsForClosedConnection() throws Exception {
connection.start();

final ActiveMQConnection activeMQConnection = (ActiveMQConnection) connection;
assertTrue("creation advisory received in time with async dispatch", Wait.waitFor(new Wait.Condition() {
@Override
public boolean isSatisified() throws Exception {
return activeMQConnection.activeTempDestinations.containsKey(queue);
}
}));
assertTrue("creation advisory received in time with async dispatch",
Wait.waitFor(() -> activeMQConnection.activeTempDestinations.contains(queue)));

// This message delivery should work since the temp connection is still
// open.
Expand Down Expand Up @@ -268,12 +264,8 @@ public void testPublishFailsForDestroyedTempDestination() throws Exception {
connection.start();

final ActiveMQConnection activeMQConnection = (ActiveMQConnection) connection;
assertTrue("creation advisory received in time with async dispatch", Wait.waitFor(new Wait.Condition() {
@Override
public boolean isSatisified() throws Exception {
return activeMQConnection.activeTempDestinations.containsKey(queue);
}
}));
assertTrue("creation advisory received in time with async dispatch",
Wait.waitFor((Wait.Condition) () -> activeMQConnection.activeTempDestinations.contains(queue)));

// This message delivery should work since the temp connection is still
// open.
Expand Down
Loading