GitOrigin-RevId: 727d5a0eb87d0b6a303c19bd80c54b3065ed4e4a
This commit is contained in:
Vladimir Krivosheev
2019-08-01 14:03:25 +03:00
committed by intellij-monorepo-bot
parent 8418f9324a
commit b350b3472f
2 changed files with 7 additions and 7 deletions
@@ -116,7 +116,7 @@ final class MessageBusConnectionImpl implements MessageBusConnection {
final Message messageOnLocalQueue = myPendingMessages.get().poll();
assert messageOnLocalQueue == message;
final Topic topic = message.getTopic();
Topic<?> topic = message.getTopic();
Object handler = mySubscriptions.get(topic);
try {
if (handler == myDefaultHandler) {
@@ -154,7 +154,7 @@ final class MessageBusConnectionImpl implements MessageBusConnection {
myPendingMessages.get().offer(message);
}
boolean containsMessage(@NotNull Topic topic) {
boolean containsMessage(@NotNull Topic<?> topic) {
Queue<Message> pendingMessages = myPendingMessages.get();
if (pendingMessages.isEmpty()) return false;
@@ -38,17 +38,17 @@ public class MessageBusImpl implements MessageBus {
*/
private final int[] myOrder;
private final ConcurrentMap<Topic, Object> myPublishers = ContainerUtil.newConcurrentMap();
private final ConcurrentMap<Topic<?>, Object> myPublishers = ContainerUtil.newConcurrentMap();
/**
* This bus's subscribers
*/
private final ConcurrentMap<Topic, List<MessageBusConnectionImpl>> mySubscribers = ContainerUtil.newConcurrentMap();
private final ConcurrentMap<Topic<?>, List<MessageBusConnectionImpl>> mySubscribers = ContainerUtil.newConcurrentMap();
/**
* Caches subscribers for this bus and its children or parent, depending on the topic's broadcast policy
*/
private final Map<Topic, List<MessageBusConnectionImpl>> mySubscriberCache = ContainerUtil.newConcurrentMap();
private final Map<Topic<?>, List<MessageBusConnectionImpl>> mySubscriberCache = ContainerUtil.newConcurrentMap();
private final List<MessageBusImpl> myChildBuses = ContainerUtil.createLockFreeCopyOnWriteList();
@NotNull
@@ -296,7 +296,7 @@ public class MessageBusImpl implements MessageBus {
return myOwner;
}
private void calcSubscribers(@NotNull Topic topic, @NotNull List<? super MessageBusConnectionImpl> result) {
private void calcSubscribers(@NotNull Topic<?> topic, @NotNull List<? super MessageBusConnectionImpl> result) {
final List<MessageBusConnectionImpl> topicSubscribers = mySubscribers.get(topic);
if (topicSubscribers != null) {
result.addAll(topicSubscribers);
@@ -328,7 +328,7 @@ public class MessageBusImpl implements MessageBus {
}
@NotNull
private List<MessageBusConnectionImpl> getTopicSubscribers(@NotNull Topic topic) {
private List<MessageBusConnectionImpl> getTopicSubscribers(@NotNull Topic<?> topic) {
List<MessageBusConnectionImpl> topicSubscribers = mySubscriberCache.get(topic);
if (topicSubscribers == null) {
topicSubscribers = new SmartList<>();