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