Skip to content

Commit e1626b6

Browse files
committed
add option to socket middleware that allows connecting to all resolved addresses for a given host.
Clean up code around connecting to InetSocketAddress: Only resolve the host if necessary. Resolve the address on a separate thread pool so it doesnt block the NIO thread. Use TransformFuture to handle the conversion from an InetAddress to a AsyncSocket.
1 parent ba2bb57 commit e1626b6

5 files changed

Lines changed: 256 additions & 112 deletions

File tree

AndroidAsync/src/com/koushikdutta/async/AsyncServer.java

Lines changed: 94 additions & 81 deletions
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,11 @@
77
import com.koushikdutta.async.callback.ConnectCallback;
88
import com.koushikdutta.async.callback.ListenCallback;
99
import com.koushikdutta.async.future.Cancellable;
10+
import com.koushikdutta.async.future.Future;
11+
import com.koushikdutta.async.future.FutureCallback;
1012
import com.koushikdutta.async.future.SimpleCancellable;
1113
import com.koushikdutta.async.future.SimpleFuture;
14+
import com.koushikdutta.async.future.TransformFuture;
1215

1316
import java.io.IOException;
1417
import java.net.InetAddress;
@@ -20,6 +23,8 @@
2023
import java.util.LinkedList;
2124
import java.util.Set;
2225
import java.util.WeakHashMap;
26+
import java.util.concurrent.ExecutorService;
27+
import java.util.concurrent.Executors;
2328
import java.util.concurrent.Semaphore;
2429
import java.util.concurrent.TimeUnit;
2530

@@ -365,111 +370,119 @@ public void stop() {
365370
}
366371
});
367372
}
368-
369-
private void connectSocketInternal(final SocketChannel socket, final SocketAddress remote, ConnectFuture cancel) {
370-
if (cancel.isCancelled())
371-
return;
372-
SelectionKey ckey = null;
373-
try {
374-
socket.configureBlocking(false);
375-
ckey = socket.register(mSelector, SelectionKey.OP_CONNECT);
376-
ckey.attach(cancel);
377-
socket.connect(remote);
378-
}
379-
catch (Exception e) {
380-
if (ckey != null)
381-
ckey.cancel();
382-
if (cancel.setComplete(e))
383-
cancel.callback.onConnectCompleted(e, null);
384-
}
385-
}
386-
373+
387374
private class ConnectFuture extends SimpleFuture<AsyncNetworkSocket> {
388375
@Override
389-
public boolean cancel() {
390-
if (!super.cancel())
391-
return false;
392-
376+
protected void cancelCleanup() {
377+
super.cancelCleanup();
393378
post(new Runnable() {
394379
@Override
395380
public void run() {
396381
try {
397-
socket.close();
382+
if (socket != null)
383+
socket.close();
398384
}
399385
catch (IOException e) {
400386
}
401387
}
402388
});
403-
return true;
404389
}
405390

406391
SocketChannel socket;
407392
ConnectCallback callback;
408393
}
409394

410-
private ConnectFuture prepareConnectSocketCancelable(final SocketChannel socket, ConnectCallback handler) {
411-
ConnectFuture cancelable = new ConnectFuture();
412-
cancelable.socket = socket;
413-
cancelable.callback = handler;
414-
return cancelable;
415-
}
416-
417-
public Cancellable connectSocket(final InetSocketAddress remote, final ConnectCallback handler) {
418-
try {
419-
final SocketChannel socket = SocketChannel.open();
420-
final ConnectFuture cancel = prepareConnectSocketCancelable(socket, handler);
421-
post(new Runnable() {
422-
@Override
423-
public void run() {
424-
connectSocketInternal(socket, new InetSocketAddress(remote.getHostName(), remote.getPort()), cancel);
425-
}
426-
});
427-
return cancel;
428-
}
429-
catch (final Exception e) {
430-
post(new Runnable() {
431-
@Override
432-
public void run() {
433-
handler.onConnectCompleted(e, null);
434-
}
435-
});
436-
return SimpleCancellable.COMPLETED;
437-
}
438-
}
439-
440-
public Cancellable connectSocket(final String host, final int port, final ConnectCallback handler) {
441-
try {
442-
final SocketChannel socket = SocketChannel.open();
443-
final ConnectFuture cancel = prepareConnectSocketCancelable(socket, handler);
395+
private ConnectFuture connectResolvedInetSocketAddress(final InetSocketAddress address, final ConnectCallback callback) {
396+
final ConnectFuture cancel = new ConnectFuture();
397+
assert !address.isUnresolved();
444398

445-
post(new Runnable() {
446-
@Override
447-
public void run() {
448-
SocketAddress remote;
399+
post(new Runnable() {
400+
@Override
401+
public void run() {
402+
if (cancel.isCancelled())
403+
return;
404+
405+
cancel.callback = callback;
406+
SelectionKey ckey = null;
407+
SocketChannel socket = null;
408+
try {
409+
socket = cancel.socket = SocketChannel.open();
410+
socket.configureBlocking(false);
411+
ckey = socket.register(mSelector, SelectionKey.OP_CONNECT);
412+
ckey.attach(cancel);
413+
socket.connect(address);
414+
}
415+
catch (Exception e) {
416+
if (ckey != null)
417+
ckey.cancel();
449418
try {
450-
remote = new InetSocketAddress(host, port);
419+
if (socket != null)
420+
socket.close();
451421
}
452-
catch (Exception e) {
453-
if (cancel.setComplete(e))
454-
handler.onConnectCompleted(e, null);
455-
return;
422+
catch (Exception ignored) {
456423
}
457-
458-
connectSocketInternal(socket, remote, cancel);
424+
cancel.setComplete(e);
459425
}
460-
});
461-
462-
return cancel;
426+
}
427+
});
428+
429+
return cancel;
430+
}
431+
432+
public Cancellable connectSocket(final InetSocketAddress remote, final ConnectCallback callback) {
433+
if (!remote.isUnresolved())
434+
return connectResolvedInetSocketAddress(remote, callback);
435+
436+
return new TransformFuture<AsyncSocket, InetAddress>() {
437+
@Override
438+
protected void transform(InetAddress result) throws Exception {
439+
setParent(connectResolvedInetSocketAddress(new InetSocketAddress(remote.getHostName(), remote.getPort()), callback));
440+
}
463441
}
464-
catch (final Exception e) {
465-
post(new Runnable() {
466-
@Override
467-
public void run() {
468-
handler.onConnectCompleted(e, null);
442+
.from(getByName(remote.getHostName()));
443+
}
444+
445+
public Cancellable connectSocket(final String host, final int port, final ConnectCallback callback) {
446+
return connectSocket(InetSocketAddress.createUnresolved(host, port), callback);
447+
}
448+
449+
ExecutorService synchronousWorkers = Executors.newFixedThreadPool(4);
450+
public Future<InetAddress[]> getAllByName(final String host) {
451+
final SimpleFuture<InetAddress[]> ret = new SimpleFuture<InetAddress[]>();
452+
synchronousWorkers.execute(new Runnable() {
453+
@Override
454+
public void run() {
455+
try {
456+
final InetAddress[] result = InetAddress.getAllByName(host);
457+
if (result == null || result.length == 0)
458+
throw new Exception("no addresses for host");
459+
post(new Runnable() {
460+
@Override
461+
public void run() {
462+
ret.setComplete(null, result);
463+
}
464+
});
465+
} catch (final Exception e) {
466+
post(new Runnable() {
467+
@Override
468+
public void run() {
469+
ret.setComplete(e, null);
470+
}
471+
});
469472
}
470-
});
471-
return SimpleCancellable.COMPLETED;
473+
}
474+
});
475+
return ret;
476+
}
477+
478+
public Future<InetAddress> getByName(String host) {
479+
return new TransformFuture<InetAddress, InetAddress[]>() {
480+
@Override
481+
protected void transform(InetAddress[] result) throws Exception {
482+
setComplete(result[0]);
483+
}
472484
}
485+
.from(getAllByName(host));
473486
}
474487

475488
public AsyncDatagramSocket connectDatagram(final String host, final int port) throws IOException {
@@ -800,7 +813,7 @@ else if (key.isConnectable()) {
800813
catch (Exception ex) {
801814
key.cancel();
802815
sc.close();
803-
if (cancel.setComplete())
816+
if (cancel.setComplete(ex))
804817
cancel.callback.onConnectCompleted(ex, null);
805818
}
806819
}

AndroidAsync/src/com/koushikdutta/async/future/TransformFuture.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@
33
public abstract class TransformFuture<T, F> extends SimpleFuture<T> implements FutureCallback<F> {
44
@Override
55
public void onCompleted(Exception e, F result) {
6+
if (isCancelled())
7+
return;
68
if (e != null) {
79
error(e);
810
return;

AndroidAsync/src/com/koushikdutta/async/http/AsyncHttpClient.java

Lines changed: 12 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -43,11 +43,21 @@ public void insertMiddleware(AsyncHttpClientMiddleware middleware) {
4343
mMiddleware.add(0, middleware);
4444
}
4545

46+
AsyncSSLSocketMiddleware sslSocketMiddleware;
47+
AsyncSocketMiddleware socketMiddleware;
4648
AsyncServer mServer;
4749
public AsyncHttpClient(AsyncServer server) {
4850
mServer = server;
49-
insertMiddleware(new AsyncSocketMiddleware(this));
50-
insertMiddleware(new AsyncSSLSocketMiddleware(this));
51+
insertMiddleware(socketMiddleware = new AsyncSocketMiddleware(this));
52+
insertMiddleware(sslSocketMiddleware = new AsyncSSLSocketMiddleware(this));
53+
}
54+
55+
public AsyncSocketMiddleware getSocketMiddleware() {
56+
return socketMiddleware;
57+
}
58+
59+
public AsyncSSLSocketMiddleware getSSLSocketMiddleware() {
60+
return sslSocketMiddleware;
5161
}
5262

5363
public Future<AsyncHttpResponse> execute(final AsyncHttpRequest request, final HttpConnectCallback callback) {

0 commit comments

Comments
 (0)