001 /*
002 GRANITE DATA SERVICES
003 Copyright (C) 2012 GRANITE DATA SERVICES S.A.S.
004
005 This file is part of Granite Data Services.
006
007 Granite Data Services is free software; you can redistribute it and/or modify
008 it under the terms of the GNU Library General Public License as published by
009 the Free Software Foundation; either version 2 of the License, or (at your
010 option) any later version.
011
012 Granite Data Services is distributed in the hope that it will be useful, but
013 WITHOUT ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or
014 FITNESS FOR A PARTICULAR PURPOSE. See the GNU Library General Public License
015 for more details.
016
017 You should have received a copy of the GNU Library General Public License
018 along with this library; if not, see <http://www.gnu.org/licenses/>.
019 */
020
021 package org.granite.client.messaging;
022
023 import java.util.concurrent.ConcurrentHashMap;
024
025 import org.granite.client.messaging.channel.MessagingChannel;
026 import org.granite.client.messaging.channel.ResponseMessageFuture;
027 import org.granite.client.messaging.events.IssueEvent;
028 import org.granite.client.messaging.events.ResultEvent;
029 import org.granite.client.messaging.events.TopicMessageEvent;
030 import org.granite.client.messaging.messages.push.TopicMessage;
031 import org.granite.client.messaging.messages.requests.SubscribeMessage;
032 import org.granite.client.messaging.messages.requests.UnsubscribeMessage;
033 import org.granite.logging.Logger;
034
035 /**
036 * @author Franck WOLFF
037 */
038 public class Consumer extends AbstractTopicAgent {
039
040 private static final Logger log = Logger.getLogger(Consumer.class);
041
042 private final ConcurrentHashMap<TopicMessageListener, Boolean> listeners = new ConcurrentHashMap<TopicMessageListener, Boolean>();
043
044 private String subscriptionId = null;
045 private String selector = null;
046
047 public Consumer(MessagingChannel channel, String destination, String topic) {
048 super(channel, destination, topic);
049 }
050
051 public String getSelector() {
052 return selector;
053 }
054
055 public void setSelector(String selector) {
056 this.selector = selector;
057 }
058
059 public boolean isSubscribed() {
060 return subscriptionId != null;
061 }
062
063 public String getSubscriptionId() {
064 return subscriptionId;
065 }
066
067 public ResponseMessageFuture subscribe(ResponseListener...listeners) {
068 SubscribeMessage subscribeMessage = new SubscribeMessage(destination, topic, selector);
069 subscribeMessage.getHeaders().putAll(defaultHeaders);
070
071 final Consumer consumer = this;
072 ResponseListener listener = new ResultIssuesResponseListener() {
073
074 @Override
075 public void onResult(ResultEvent event) {
076 subscriptionId = (String)event.getResult();
077 channel.addConsumer(consumer);
078 }
079
080 @Override
081 public void onIssue(IssueEvent event) {
082 log.error("Subscription failed %s: %s", consumer, event);
083 }
084 };
085
086 if (listeners == null || listeners.length == 0)
087 listeners = new ResponseListener[]{listener};
088 else {
089 ResponseListener[] tmp = new ResponseListener[listeners.length + 1];
090 System.arraycopy(listeners, 0, tmp, 0, listeners.length);
091 tmp[listeners.length] = listener;
092 listeners = tmp;
093 }
094
095 return channel.send(subscribeMessage, listeners);
096 }
097
098 public ResponseMessageFuture unsubscribe(ResponseListener...listeners) {
099 UnsubscribeMessage unsubscribeMessage = new UnsubscribeMessage(destination, topic, subscriptionId);
100 unsubscribeMessage.getHeaders().putAll(defaultHeaders);
101
102 final Consumer consumer = this;
103 ResponseListener listener = new ResultIssuesResponseListener() {
104
105 @Override
106 public void onResult(ResultEvent event) {
107 channel.removeConsumer(consumer);
108 subscriptionId = null;
109 }
110
111 @Override
112 public void onIssue(IssueEvent event) {
113 log.error("Unsubscription failed %s: %s", consumer, event);
114 }
115 };
116
117 if (listeners == null || listeners.length == 0)
118 listeners = new ResponseListener[]{listener};
119 else {
120 ResponseListener[] tmp = new ResponseListener[listeners.length + 1];
121 System.arraycopy(listeners, 0, tmp, 0, listeners.length);
122 tmp[listeners.length] = listener;
123 listeners = tmp;
124 }
125
126 return channel.send(unsubscribeMessage, listeners);
127 }
128
129 public void addMessageListener(TopicMessageListener listener) {
130 listeners.putIfAbsent(listener, Boolean.TRUE);
131 }
132
133 public boolean removeMessageListener(TopicMessageListener listener) {
134 return listeners.remove(listener) != null;
135 }
136
137 public void onDisconnect() {
138 subscriptionId = null;
139 }
140
141 public void onMessage(TopicMessage message) {
142 for (TopicMessageListener listener : listeners.keySet()) {
143 try {
144 listener.onMessage(new TopicMessageEvent(this, message));
145 }
146 catch (Exception e) {
147 log.error(e, "Consumer listener threw an exception: ", listener);
148 }
149 }
150 }
151
152 @Override
153 public String toString() {
154 return getClass().getName() + " {subscriptionId=" + subscriptionId +
155 ", destination=" + destination +
156 ", topic=" + topic +
157 ", selector=" + selector +
158 "}";
159 }
160 }