diff --git a/client/src/main/java/org/asynchttpclient/netty/request/NettyRequestSender.java b/client/src/main/java/org/asynchttpclient/netty/request/NettyRequestSender.java index 06a7fc444..abf5fba63 100755 --- a/client/src/main/java/org/asynchttpclient/netty/request/NettyRequestSender.java +++ b/client/src/main/java/org/asynchttpclient/netty/request/NettyRequestSender.java @@ -18,7 +18,6 @@ import io.netty.bootstrap.Bootstrap; import io.netty.buffer.ByteBuf; import io.netty.channel.Channel; -import io.netty.channel.ChannelFuture; import io.netty.channel.ChannelInitializer; import io.netty.channel.ChannelProgressivePromise; import io.netty.channel.ChannelPromise; @@ -736,17 +735,17 @@ public void writeRequest(NettyResponseFuture future, Channel channel) { return; } - // if the request has a body, we want to track progress + // Listen before writing: off the event loop the write can complete first, and a listener added to + // a completed promise is then notified after the response has been read. if (writeBody) { // FIXME does this really work??? the promise is for the request without body!!! ChannelProgressivePromise promise = channel.newProgressivePromise(); - ChannelFuture f = channel.write(httpRequest, promise); - f.addListener(new WriteProgressListener(future, true, 0L)); + promise.addListener(new WriteProgressListener(future, true, 0L)); + channel.write(httpRequest, promise); } else { - // we can just track write completion ChannelPromise promise = channel.newPromise(); - ChannelFuture f = channel.writeAndFlush(httpRequest, promise); - f.addListener(new WriteCompleteListener(future)); + promise.addListener(new WriteCompleteListener(future)); + channel.writeAndFlush(httpRequest, promise); } } diff --git a/client/src/test/java/org/asynchttpclient/channel/ConnectionPoolTest.java b/client/src/test/java/org/asynchttpclient/channel/ConnectionPoolTest.java index 89cc5b4c2..e6fb48391 100644 --- a/client/src/test/java/org/asynchttpclient/channel/ConnectionPoolTest.java +++ b/client/src/test/java/org/asynchttpclient/channel/ConnectionPoolTest.java @@ -245,6 +245,27 @@ public void nonPoolableConnectionReleaseSemaphoresTest() throws Throwable { } } + @Test + public void headersWrittenIsReportedBeforeCompletionOnAPooledConnection() throws Exception { + RequestBuilder request = get("http://localhost:" + port1 + "/Test"); + + // The race hits roughly one request in ten, so a few hundred make a miss vanishingly unlikely. + int outOfOrder = 0; + try (AsyncHttpClient client = asyncHttpClient()) { + for (int i = 0; i < 500; i++) { + EventCollectingHandler handler = new EventCollectingHandler(); + client.executeRequest(request, handler).get(3, TimeUnit.SECONDS); + handler.waitForCompletion(3, TimeUnit.SECONDS); + List events = new ArrayList<>(handler.firedEvents); + int written = events.indexOf(HEADERS_WRITTEN_EVENT); + if (written < 0 || written > events.indexOf(COMPLETED_EVENT)) { + outOfOrder++; + } + } + } + assertEquals(0, outOfOrder, "requests whose headers-written event came after completion"); + } + @Test public void testPooledEventsFired() throws Exception { RequestBuilder request = get("http://localhost:" + port1 + "/Test"); @@ -261,7 +282,7 @@ public void testPooledEventsFired() throws Exception { Object[] expectedEvents = {CONNECTION_POOL_EVENT, CONNECTION_POOLED_EVENT, REQUEST_SEND_EVENT, HEADERS_WRITTEN_EVENT, STATUS_RECEIVED_EVENT, HEADERS_RECEIVED_EVENT, CONNECTION_OFFER_EVENT, COMPLETED_EVENT}; - assertArrayEquals(secondHandler.firedEvents.toArray(), expectedEvents, "Got " + Arrays.toString(secondHandler.firedEvents.toArray())); + assertArrayEquals(expectedEvents, secondHandler.firedEvents.toArray(), "Got " + Arrays.toString(secondHandler.firedEvents.toArray())); } } }