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 }