Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -736,17 +735,17 @@ public <T> void writeRequest(NettyResponseFuture<T> 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);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<String> 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");
Expand All @@ -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()));
}
}
}
Loading