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.jetty;
022    
023    import org.granite.gravity.AbstractChannel;
024    import org.granite.gravity.AsyncHttpContext;
025    import org.granite.gravity.Gravity;
026    import org.granite.gravity.GravityConfig;
027    import org.granite.gravity.MessageReceivingException;
028    import org.granite.logging.Logger;
029    import org.mortbay.util.ajax.Continuation;
030    
031    import flex.messaging.messages.AsyncMessage;
032    
033    /**
034     * @author William DRAI
035     */
036    public class ContinuationChannel extends AbstractChannel {
037    
038        private static final Logger log = Logger.getLogger(ContinuationChannel.class);
039    
040        private Continuation continuation = null;
041    
042        public ContinuationChannel(Gravity gravity, String id, ContinuationChannelFactory factory, String clientType) {
043            super(gravity, id, factory, clientType);
044        }
045    
046        public void setContinuation(Continuation continuation) {
047            log.debug("Setting continuation %s for client: %s", continuation, getId());
048            try {
049                if (this.continuation != null && this.continuation.isPending()) {
050                    log.debug("Resuming previous continuation %s for client: %s", this.continuation, getId());
051                    this.continuation.resume();
052                }
053            }
054            finally {
055                this.continuation = continuation;
056            }
057        }
058        
059        public void close() {
060            try {
061                if (this.continuation != null)
062                    this.continuation.reset();
063            }
064            finally {
065                this.continuation = null;
066            }
067        }
068        
069    
070        public void resume() {
071            log.debug("Resuming pending continuation %s for client: %s", continuation, getId());
072            try {
073                if (this.continuation != null && this.continuation.isPending())
074                    this.continuation.resume();
075            }
076            finally {
077                this.continuation = null;
078            }
079        }
080        
081        @Override
082            public void receive(AsyncMessage message) throws MessageReceivingException {
083                    if (message == null)
084                            throw new NullPointerException("message cannot be null");
085    
086                    GravityConfig gravityConfig = gravity.getGravityConfig();
087                    
088                    receivedQueueLock.lock();
089                    try {
090                            if (receivedQueue.size() + 1 > gravityConfig.getMaxMessagesQueuedPerChannel())
091                                    throw new MessageReceivingException(message, "Could not queue message (channel's queue is full) for channel: " + this);
092                            
093                            receivedQueue.add(message);
094                    }
095                    finally {
096                            receivedQueueLock.unlock();
097                    }
098                    
099                    synchronized (this) {
100                            resume();
101                    }
102            }
103    
104            @Override
105            protected boolean hasAsyncHttpContext() {
106                    return false;
107            }
108    
109            @Override
110            protected void releaseAsyncHttpContext(AsyncHttpContext context) {
111            }
112    
113            @Override
114            protected AsyncHttpContext acquireAsyncHttpContext() {
115            return null;
116        }
117    
118        public boolean isLocal() {
119            return true;
120        }
121    
122            @Override
123            public void destroy() {
124                    try {
125                            super.destroy();
126                    }
127                    finally {
128                            synchronized (this) {
129                                    close();
130                            }
131                    }
132            }
133    }