This commit is contained in:
Alexey Kudravtsev
2012-04-10 12:05:51 +04:00
parent b70c17f1cf
commit efb3edebb2
2 changed files with 21 additions and 7 deletions
@@ -47,6 +47,7 @@ public class MessageBusConnectionImpl implements MessageBusConnection {
myBus = bus;
}
@Override
public <L> void subscribe(Topic<L> topic, L handler) throws IllegalStateException {
if (mySubscriptions.put(topic, handler) != null) {
throw new IllegalStateException("Subscription to " + topic + " already exists");
@@ -54,6 +55,7 @@ public class MessageBusConnectionImpl implements MessageBusConnection {
myBus.notifyOnSubscription(this, topic);
}
@Override
@SuppressWarnings("unchecked")
public <L> void subscribe(Topic<L> topic) throws IllegalStateException {
if (myDefaultHandler == null) {
@@ -61,32 +63,34 @@ public class MessageBusConnectionImpl implements MessageBusConnection {
+ "Target topic: " + topic);
}
if (topic.getListenerClass().isInstance(myDefaultHandler)) {
throw new IllegalStateException(String.format(
"Can't subscribe to the topic '%s'. Reason: default handler has incompatible type - expected: '%s', actual: '%s'",
topic, topic.getListenerClass(), myDefaultHandler.getClass()
));
throw new IllegalStateException("Can't subscribe to the topic '" + topic +"'. Default handler has incompatible type - expected: '" +
topic.getListenerClass() + "', actual: '" + myDefaultHandler.getClass() + "'");
}
subscribe(topic, (L)myDefaultHandler);
}
@Override
public void setDefaultHandler(MessageHandler handler) {
myDefaultHandler = handler;
}
@Override
public void disconnect() {
Queue<Message> jobs = myPendingMessages.get();
if (!jobs.isEmpty()) {
LOG.error("Not delivered events in the queue: "+jobs);
}
myPendingMessages.remove();
myBus.notifyConnectionTerminated(this);
if (!jobs.isEmpty()) {
LOG.error("Not delivered events in the queue: " + jobs);
}
}
@Override
public void dispose() {
disconnect();
}
@Override
public void deliverImmediately() {
while (!myPendingMessages.get().isEmpty()) {
myBus.deliverSingleMessage();
@@ -27,6 +27,7 @@ import com.intellij.util.containers.ContainerUtil;
import com.intellij.util.messages.MessageBus;
import com.intellij.util.messages.MessageBusConnection;
import com.intellij.util.messages.Topic;
import org.jetbrains.annotations.NonNls;
import org.jetbrains.annotations.NotNull;
import java.lang.reflect.InvocationHandler;
@@ -60,6 +61,7 @@ public class MessageBusImpl implements MessageBus {
private final Object myOwner;
private boolean myDisposed;
@SuppressWarnings("UnusedDeclaration")
public MessageBusImpl() {
this(null, null);
}
@@ -97,18 +99,21 @@ public class MessageBusImpl implements MessageBus {
public final MessageBusConnectionImpl connection;
public final Message message;
@NonNls
@Override
public String toString() {
return "{ DJob connection:" + connection.toString() + "; message: " + message + " }";
}
}
@Override
@NotNull
public MessageBusConnection connect() {
checkNotDisposed();
return new MessageBusConnectionImpl(this);
}
@Override
@NotNull
public MessageBusConnection connect(@NotNull Disposable parentDisposable) {
final MessageBusConnection connection = connect();
@@ -116,6 +121,7 @@ public class MessageBusImpl implements MessageBus {
return connection;
}
@Override
@NotNull
@SuppressWarnings({"unchecked"})
public <L> L syncPublisher(@NotNull final Topic<L> topic) {
@@ -124,6 +130,7 @@ public class MessageBusImpl implements MessageBus {
if (publisher == null) {
final Class<L> listenerClass = topic.getListenerClass();
InvocationHandler handler = new InvocationHandler() {
@Override
public Object invoke(Object proxy, Method method, Object[] args) throws Throwable {
sendMessage(new Message(topic, method, args));
return NA;
@@ -135,6 +142,7 @@ public class MessageBusImpl implements MessageBus {
return publisher;
}
@Override
@NotNull
@SuppressWarnings({"unchecked"})
public <L> L asyncPublisher(@NotNull final Topic<L> topic) {
@@ -143,6 +151,7 @@ public class MessageBusImpl implements MessageBus {
if (publisher == null) {
final Class<L> listenerClass = topic.getListenerClass();
InvocationHandler handler = new InvocationHandler() {
@Override
public Object invoke(Object proxy, Method method, Object[] args) throws Throwable {
postMessage(new Message(topic, method, args));
return NA;
@@ -154,6 +163,7 @@ public class MessageBusImpl implements MessageBus {
return publisher;
}
@Override
public void dispose() {
checkNotDisposed();
Queue<DeliveryJob> jobs = myMessageQueue.get();