001    /*
002      GRANITE DATA SERVICES
003      Copyright (C) 2011 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.gravity.adapters;
022    
023    import java.util.concurrent.ConcurrentHashMap;
024    
025    import org.granite.gravity.AsyncPublishedMessage;
026    import org.granite.gravity.Channel;
027    import org.granite.gravity.MessagePublishingException;
028    import org.granite.logging.Logger;
029    import org.granite.messaging.service.ServiceException;
030    import org.granite.util.XMap;
031    
032    import flex.messaging.messages.AcknowledgeMessage;
033    import flex.messaging.messages.AsyncMessage;
034    import flex.messaging.messages.CommandMessage;
035    import flex.messaging.messages.ErrorMessage;
036    
037    /**
038     * @author William DRAI
039     */
040    public class SimpleServiceAdapter extends ServiceAdapter {
041    
042        private static final Logger log = Logger.getLogger(SimpleServiceAdapter.class);
043    
044        private final Topic rootTopic = new Topic("/", this);
045        private transient ConcurrentHashMap<String, TopicId> _topicIdCache;
046        
047        private boolean noLocal = false;
048    
049        @Override
050        public void configure(XMap adapterProperties, XMap destinationProperties) throws ServiceException {
051            super.configure(adapterProperties, destinationProperties);
052            
053            _topicIdCache = new ConcurrentHashMap<String, TopicId>();
054            
055            if (Boolean.TRUE.toString().equals(destinationProperties.get("no-local")))
056                    noLocal = true;
057        }
058    
059    
060        public Topic getTopic(TopicId id) {
061            return rootTopic.getChild(id);
062        }
063    
064        public Topic getTopic(String id) {
065            TopicId cid = getTopicId(id);
066            if (cid.depth() == 0)
067                return null;
068            return rootTopic.getChild(cid);
069        }
070    
071        public Topic getTopic(String id, boolean create)  {
072            synchronized (this) {
073                Topic topic = getTopic(id);
074    
075                if (topic == null && create) {
076                    topic = new Topic(id, this);
077                    rootTopic.addChild(topic);
078                    log.debug("New Topic: %s", topic);
079                }
080                return topic;
081            }
082        }
083    
084        public TopicId getTopicId(String id) {
085            TopicId tid = _topicIdCache.get(id);
086            if (tid == null) {
087                tid = new TopicId(id);
088                TopicId tmpTid = _topicIdCache.putIfAbsent(id, tid); 
089                if (tmpTid != null) 
090                    tid = tmpTid; 
091            }
092            return tid;
093        }
094    
095        public boolean hasTopic(String id) {
096            TopicId cid = getTopicId(id);
097            return rootTopic.getChild(cid) != null;
098        }
099    
100        @Override
101        public Object invoke(Channel fromChannel, AsyncMessage message) {
102            String topicId = TopicId.normalize(((String)message.getHeader(AsyncMessage.SUBTOPIC_HEADER)));
103    
104            AsyncMessage reply = null;
105    
106            if (getSecurityPolicy().canPublish(fromChannel, topicId, message)) {
107                TopicId tid = getTopicId(topicId);
108    
109                    try {
110                                    fromChannel.publish(new AsyncPublishedMessage(rootTopic, tid, message));
111                        reply = new AcknowledgeMessage(message);
112                        reply.setMessageId(message.getMessageId());
113                            }
114                    catch (MessagePublishingException e) {
115                                    log.error(e, "Error while publishing message: %s from channel %s to topic: %s", message, fromChannel, tid);
116                        reply = new ErrorMessage(message, null);
117                        ((ErrorMessage)reply).setFaultString("Server.Publish.Error");
118                            }
119            }
120            else {
121                    log.warn("Channel %s tried to publish a message to topic %s", fromChannel, topicId);
122                reply = new ErrorMessage(message, null);
123                ((ErrorMessage)reply).setFaultString("Server.Publish.Denied");
124            }
125    
126            return reply;
127        }
128    
129        @Override
130        public Object manage(Channel fromChannel, CommandMessage message) {
131            AsyncMessage reply = null;
132    
133            if (message.getOperation() == CommandMessage.SUBSCRIBE_OPERATION) {
134                String subscribeTopicId = TopicId.normalize(((String)message.getHeader(AsyncMessage.SUBTOPIC_HEADER)));
135    
136                if (getSecurityPolicy().canSubscribe(fromChannel, subscribeTopicId, message)) {
137                    Topic topic = getTopic(subscribeTopicId);
138                    if (topic == null && getSecurityPolicy().canCreate(fromChannel, subscribeTopicId, message))
139                        topic = getTopic(subscribeTopicId, true);
140    
141                    if (topic != null) {
142                        String subscriptionId = (String)message.getHeader(AsyncMessage.DESTINATION_CLIENT_ID_HEADER);
143                        String selector = (String)message.getHeader(CommandMessage.SELECTOR_HEADER);
144                        if (subscriptionId == null)
145                            log.warn("No subscriptionId for subscription message");
146                        else
147                            topic.subscribe(fromChannel, message.getDestination(), subscriptionId, selector, noLocal);
148    
149                        reply = new AcknowledgeMessage(message);
150                    }
151                    else {
152                        reply = new ErrorMessage(message, null);
153                        ((ErrorMessage)reply).setFaultString("Server.CreateTopic.Denied");
154                    }
155                }
156                else {
157                    reply = new ErrorMessage(message, null);
158                    ((ErrorMessage)reply).setFaultString("Server.Subscribe.Denied");
159                }
160            }
161            else if (message.getOperation() == CommandMessage.UNSUBSCRIBE_OPERATION) {
162                String unsubscribeTopicId = TopicId.normalize(((String)message.getHeader(AsyncMessage.SUBTOPIC_HEADER)));
163    
164                Topic topic = getTopic(unsubscribeTopicId);
165                String subscriptionId = null;
166                if (topic != null) {
167                    subscriptionId = (String)message.getHeader(AsyncMessage.DESTINATION_CLIENT_ID_HEADER);
168                    if (subscriptionId == null)
169                            log.warn("No subscriptionId for unsubscription message");
170                    else
171                            topic.unsubscribe(fromChannel, subscriptionId);
172                }
173    
174                reply = new AcknowledgeMessage(message);
175                reply.setHeader(AsyncMessage.DESTINATION_CLIENT_ID_HEADER, subscriptionId);
176            }
177            else {
178                reply = new ErrorMessage(message, null);
179                ((ErrorMessage)reply).setFaultString("unknown operation");
180    
181            }
182    
183            return reply;
184        }
185    }