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.channel;
022    
023    import java.io.IOException;
024    import java.io.InputStream;
025    import java.io.UnsupportedEncodingException;
026    import java.net.URI;
027    import java.util.Timer;
028    import java.util.TimerTask;
029    import java.util.concurrent.BlockingQueue;
030    import java.util.concurrent.ConcurrentHashMap;
031    import java.util.concurrent.ConcurrentMap;
032    import java.util.concurrent.ExecutionException;
033    import java.util.concurrent.LinkedBlockingQueue;
034    import java.util.concurrent.Semaphore;
035    import java.util.concurrent.TimeUnit;
036    import java.util.concurrent.TimeoutException;
037    
038    import org.granite.client.messaging.AllInOneResponseListener;
039    import org.granite.client.messaging.ResponseListener;
040    import org.granite.client.messaging.ResponseListenerDispatcher;
041    import org.granite.client.messaging.ResultFaultIssuesResponseListener;
042    import org.granite.client.messaging.events.Event;
043    import org.granite.client.messaging.events.FaultEvent;
044    import org.granite.client.messaging.events.IssueEvent;
045    import org.granite.client.messaging.events.ResultEvent;
046    import org.granite.client.messaging.events.Event.Type;
047    import org.granite.client.messaging.messages.MessageChain;
048    import org.granite.client.messaging.messages.RequestMessage;
049    import org.granite.client.messaging.messages.ResponseMessage;
050    import org.granite.client.messaging.messages.requests.LoginMessage;
051    import org.granite.client.messaging.messages.requests.LogoutMessage;
052    import org.granite.client.messaging.messages.requests.PingMessage;
053    import org.granite.client.messaging.messages.responses.FaultMessage;
054    import org.granite.client.messaging.messages.responses.ResultMessage;
055    //import org.granite.client.messaging.transport.HTTPTransport;
056    import org.granite.client.messaging.transport.Transport;
057    import org.granite.client.messaging.transport.TransportFuture;
058    import org.granite.client.messaging.transport.TransportMessage;
059    import org.granite.client.messaging.transport.TransportStopListener;
060    import org.granite.logging.Logger;
061    
062    /**
063     * @author Franck WOLFF
064     */
065    public abstract class AbstractHTTPChannel extends AbstractChannel<Transport> implements TransportStopListener, Runnable {
066            
067            private static final Logger log = Logger.getLogger(AbstractHTTPChannel.class);
068            
069            private final BlockingQueue<AsyncToken> tokensQueue = new LinkedBlockingQueue<AsyncToken>();
070            private final ConcurrentMap<String, AsyncToken> tokensMap = new ConcurrentHashMap<String, AsyncToken>();
071    
072            private Thread senderThread = null;
073            private Semaphore connections;
074            private Timer timer = null;
075            
076            protected volatile boolean pinged = false;
077            protected volatile boolean authenticated = false;
078            protected volatile int maxConcurrentRequests;
079            protected volatile long defaultTimeToLive = TimeUnit.MINUTES.toMillis(1L); // 1 mn.
080            
081            public AbstractHTTPChannel(Transport transport, String id, URI uri) {
082                    this(transport, id, uri, 5);
083            }
084            
085            public AbstractHTTPChannel(Transport transport, String id, URI uri, int maxConcurrentRequests) {
086                    super(transport, id, uri);
087                    
088                    if (maxConcurrentRequests < 1)
089                            throw new IllegalArgumentException("maxConcurrentRequests must be greater or equal to 1");
090                    
091                    this.maxConcurrentRequests = maxConcurrentRequests;
092            }
093            
094            protected abstract TransportMessage createTransportMessage(AsyncToken token) throws UnsupportedEncodingException;
095            
096            protected abstract ResponseMessage decodeResponse(InputStream is) throws IOException;
097    
098            protected boolean schedule(TimerTask timerTask, long delay) {
099                    if (timer != null) {
100                            timer.schedule(timerTask, delay);
101                            return true;
102                    }
103                    return false;
104            }
105            
106            public boolean isAuthenticated() {
107                    return authenticated;
108            }
109    
110            public int getMaxConcurrentRequests() {
111                    return maxConcurrentRequests;
112            }
113    
114            @Override
115            public void onStop(Transport transport) {
116                    stop();
117            }
118    
119            @Override
120            public synchronized boolean start() {
121                    if (senderThread == null) {
122                            log.info("Starting channel %s...", id);
123                            senderThread = new Thread(this);
124                            try {
125                                    timer = new Timer(id + "_timer", true);
126                                    connections = new Semaphore(maxConcurrentRequests);
127                                    senderThread.start();
128                                    
129                                    transport.addStopListener(this);
130                                    
131                                    log.info("Channel %s started.", id);
132                            }
133                            catch (Exception e) {
134                                    if (timer != null) {
135                                            timer.cancel();
136                                            timer = null;
137                                    }
138                                    connections = null;
139                                    senderThread = null;
140                                    log.error(e, "Channel %s failed to start.", id);
141                                    return false;
142                            }
143                    }
144                    return true;
145            }
146    
147            @Override
148            public synchronized boolean isStarted() {
149                    return senderThread != null;
150            }
151    
152            @Override
153            public synchronized boolean stop() {
154                    if (senderThread != null) {
155                            log.info("Stopping channel %s...", id);
156                            
157                            if (timer != null) {
158                                    try {
159                                            timer.cancel();
160                                    }
161                                    catch (Exception e) {
162                                            log.error(e, "Channel %s timer failed to stop.", id);
163                                    }
164                                    finally {
165                                            timer = null;
166                                    }
167                            }
168                            
169                            connections = null;
170                            
171                            tokensMap.clear();
172                            tokensQueue.clear();
173                            
174                            Thread thread = this.senderThread;
175                            senderThread = null;
176                            thread.interrupt();
177                            
178                            pinged = false;
179                            clientId = null;
180                            authenticated = false;
181                            
182                            return true;
183                    }
184                    return false;
185            }
186    
187            @Override
188            public void run() {
189    
190                    while (!Thread.interrupted()) {
191                            try {
192                                    AsyncToken token = tokensQueue.take();
193                                    
194                                    if (token.isDone())
195                                            continue;
196    
197                                    if (!pinged) {
198                                            ResultMessage result = sendBlockingToken(new PingMessage(clientId), token);
199                                            if (result == null)
200                                                    continue;
201                                            clientId = result.getClientId();
202                                            pinged = true;
203                                    }
204    
205                                    if (!authenticated) {
206                                            Credentials credentials = this.credentials;
207                                            if (credentials != null) {
208                                                    ResultMessage result = sendBlockingToken(new LoginMessage(clientId, credentials), token);
209                                                    if (result == null)
210                                                            continue;
211                                                    authenticated = true;
212                                            }
213                                    }
214                                    
215                                    sendToken(token);
216                            }
217                            catch (InterruptedException e) {
218                                    log.info("Channel %s stopped.", id);
219                                    break;
220                            }
221                            catch (Exception e) {
222                                    log.error(e, "Channel %s got an unexepected exception.", id);
223                            }
224                    }
225            }
226            
227            private ResultMessage sendBlockingToken(RequestMessage request, AsyncToken dependentToken) {
228                    
229                    // Make this blocking request share the timeout/timeToLive values of the dependent token.
230                    request.setTimestamp(dependentToken.getRequest().getTimestamp());
231                    request.setTimeToLive(dependentToken.getRequest().getTimeToLive());
232                    
233                    // Create the blocking token and schedule it with the dependent token timeout.
234                    AsyncToken blockingToken = new AsyncToken(request);
235                    try {
236                            timer.schedule(blockingToken, blockingToken.getRequest().getRemainingTimeToLive());
237                    }
238                    catch (IllegalArgumentException e) {
239                            dependentToken.dispatchTimeout(System.currentTimeMillis());
240                            return null;
241                    }
242                    catch (Exception e) {
243                            dependentToken.dispatchFailure(e);
244                            return null;
245                    }
246                    
247                    // Try to send the blocking token (can block if the connections semaphore can't be acquired
248                    // immediately).
249                    try {
250                            if (!sendToken(blockingToken))
251                                    return null;
252                    }
253                    catch (Exception e) {
254                            dependentToken.dispatchFailure(e);
255                            return null;
256                    }
257                    
258                    // Block until we get a server response (result or fault), a cancellation (unlikely), a timeout
259                    // or any other execution exception.
260                    try {
261                            ResponseMessage response = blockingToken.get();
262                            
263                            // Request was successful, return a non-null result. 
264                            if (response instanceof ResultMessage)
265                                    return (ResultMessage)response;
266                            
267                            if (response instanceof FaultMessage) {
268                                    FaultMessage faultMessage = (FaultMessage)response.copy(dependentToken.getRequest().getId());
269                                    if (dependentToken.getRequest() instanceof MessageChain) {
270                                            ResponseMessage nextResponse = faultMessage;
271                                            for (MessageChain<?> nextRequest = ((MessageChain<?>)dependentToken.getRequest()).getNext(); nextRequest != null; nextRequest = nextRequest.getNext()) {
272                                                    nextResponse.setNext(response.copy(nextRequest.getId()));
273                                                    nextResponse = nextResponse.getNext();
274                                            }
275                                    }
276                                    dependentToken.dispatchFault(faultMessage);
277                            }
278                            else
279                                    throw new RuntimeException("Unknow response message type: " + response);
280                            
281                    }
282                    catch (InterruptedException e) {
283                            dependentToken.dispatchFailure(e);
284                    }
285                    catch (TimeoutException e) {
286                            dependentToken.dispatchTimeout(System.currentTimeMillis());
287                    }
288                    catch (ExecutionException e) {
289                            if (e.getCause() instanceof Exception)
290                                    dependentToken.dispatchFailure((Exception)e.getCause());
291                            else
292                                    dependentToken.dispatchFailure(e);
293                    }
294                    catch (Exception e) {
295                            dependentToken.dispatchFailure(e);
296                    }
297                    
298                    return null;
299            }
300            
301            private boolean sendToken(final AsyncToken token) {
302    
303                    boolean releaseConnections = false;
304                    try {
305                        // Block until a connection is available.
306                            if (!connections.tryAcquire(token.getRequest().getRemainingTimeToLive(), TimeUnit.MILLISECONDS)) {
307                                    token.dispatchTimeout(System.currentTimeMillis());
308                                    return false;
309                            }
310    
311                            // Semaphore was successfully acquired, we must release it in the finally block unless we succeed in
312                            // sending the data (see below).
313                            releaseConnections = true;
314    
315                        // Check if the token has already received an event (likely a timeout or a cancellation).
316                            if (token.isDone())
317                                    return false;
318    
319                            // Make sure we have set a clientId (can be null for ping message).
320                            token.getRequest().setClientId(clientId);
321                            
322                        // Add the token to active tokens map.
323                        if (tokensMap.putIfAbsent(token.getId(), token) != null)
324                                    throw new RuntimeException("MessageId isn't unique: " + token.getId());
325    
326                    // Actually send the message content.
327                        TransportFuture transportFuture = transport.send(this, createTransportMessage(token));
328                        
329                        // Create and try to set a channel listener: if no event has been dispatched for this token (tokenEvent == null),
330                        // the listener will be called on the next event. Otherwise, we just call the listener immediately.
331                        ResponseListener channelListener = new ChannelResponseListener(token.getId(), tokensMap, transportFuture, connections);
332                        Event tokenEvent = token.setChannelListener(channelListener);
333                        if (tokenEvent != null)
334                                    ResponseListenerDispatcher.dispatch(channelListener, tokenEvent);
335                        
336                        // Message was sent and we were able to handle everything ourself.
337                        releaseConnections = false;
338                            
339                        return true;
340                    }
341                    catch (Exception e) {
342                            tokensMap.remove(token.getId());
343                            token.dispatchFailure(e);                       
344                            if (timer != null)
345                                    timer.purge();  // Must purge to cleanup timer references to AsyncToken
346                            return false;
347                    }
348                    finally {
349                            if (releaseConnections)
350                                    connections.release();
351                    }
352            }
353            
354            protected RequestMessage getRequest(String id) {
355                    AsyncToken token = tokensMap.get(id);
356                    return (token != null ? token.getRequest() : null);
357            }
358            
359            
360            @Override
361            public ResponseMessageFuture send(RequestMessage request, ResponseListener... listeners) {
362                    if (request == null)
363                            throw new NullPointerException("request cannot be null");
364                    
365                    if (!start())
366                            throw new RuntimeException("Channel not started");
367                    
368                    AsyncToken token = new AsyncToken(request, listeners);
369    
370                    request.setTimestamp(System.currentTimeMillis());
371                    if (request.getTimeToLive() <= 0L)
372                            request.setTimeToLive(defaultTimeToLive);
373                    
374                    try {
375                            timer.schedule(token, request.getRemainingTimeToLive());
376                            tokensQueue.add(token);
377                    }
378                    catch (Exception e) {
379                            log.error(e, "Could not add token to queue: %s", token);
380                            token.dispatchFailure(e);
381                            return new ImmediateFailureResponseMessageFuture(e);
382                    }
383                    
384                    return token;
385            }
386            
387            @Override
388            public ResponseMessageFuture logout(ResponseListener... listeners) {
389                    ResponseListener[] ls = new ResponseListener[listeners.length+1];
390                    ls[0] = new ResultFaultIssuesResponseListener() {
391                            @Override
392                            public void onResult(ResultEvent event) {
393                                    authenticated = false;
394                            }
395                            
396                            @Override
397                            public void onFault(FaultEvent event) {
398                                    authenticated = false;
399                            }
400                            
401                            @Override
402                            public void onIssue(IssueEvent event) {
403                            }
404                    };
405                    for (int i = 0; i < listeners.length; i++)
406                            ls[i+1] = listeners[i];
407                    return send(new LogoutMessage(), ls);
408            }
409    
410            @Override
411            public void onMessage(InputStream is) {
412                    try {
413                            ResponseMessage response = decodeResponse(is);
414                            
415                            if (response != null) {
416                                    
417                                    AsyncToken token = tokensMap.remove(response.getCorrelationId());
418                                    if (token == null) {
419                                            log.warn("Unknown correlation id: %s", response.getCorrelationId());
420                                            return;
421                                    }
422    
423                                    switch (response.getType()) {
424                                            case RESULT:
425                                                    token.dispatchResult((ResultMessage)response);
426                                                    break;
427                                            case FAULT:
428                                                    token.dispatchFault((FaultMessage)response);
429                                                    break;
430                                            default:
431                                                    token.dispatchFailure(new RuntimeException("Unknown message type: " + response));
432                                                    break;
433                                    }
434                                    
435                                    if (timer != null)
436                                            timer.purge();  // Must purge to cleanup timer references to AsyncToken
437                            }
438                    }
439                    catch (Exception e) {
440                            log.error(e, "Could not deserialize or dispatch incoming messages");
441                    }
442            }
443    
444            @Override
445            public void onError(TransportMessage message, Exception e) {
446                    if (message != null) {
447                            AsyncToken token = tokensMap.remove(message.getId());
448                            if (token != null) {
449                                    token.dispatchFailure(e);
450                                    if (timer != null)
451                                            timer.purge();  // Must purge to cleanup timer references to AsyncToken
452                            }
453                    }
454            }
455    
456            @Override
457            public void onCancelled(TransportMessage message) {
458                    AsyncToken token = tokensMap.remove(message.getId());
459                    if (token != null) {
460                            token.dispatchCancelled();
461                            if (timer != null)
462                                    timer.purge();  // Must purge to cleanup timer references to AsyncToken
463                    }
464            }
465            
466            private static class ChannelResponseListener extends AllInOneResponseListener {
467                    
468                    private final String tokenId;
469                    private final ConcurrentMap<String, AsyncToken> tokensMap;
470                    private final TransportFuture transportFuture;
471                    private final Semaphore connections;
472                    
473                    public ChannelResponseListener(
474                            String tokenId,
475                            ConcurrentMap<String, AsyncToken> tokensMap,
476                            TransportFuture transportFuture,
477                            Semaphore connections) {
478    
479                            this.tokenId = tokenId;
480                            this.tokensMap = tokensMap;
481                            this.transportFuture = transportFuture;
482                            this.connections = connections;
483                    }
484    
485                    @Override
486                    public void onEvent(Event event) {
487                            try {
488                                    tokensMap.remove(tokenId);
489                                    if (event.getType() == Type.TIMEOUT || event.getType() == Type.CANCELLED) {
490                                            if (transportFuture != null)
491                                                    transportFuture.cancel();
492                                    }
493                            }
494                            finally {
495                                    connections.release();
496                            }
497                    }
498            } 
499    }