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    }