diff --git a/vertx-core/src/main/java/io/vertx/core/http/impl/http1/Http1ServerConnection.java b/vertx-core/src/main/java/io/vertx/core/http/impl/http1/Http1ServerConnection.java index 2b108efe9a1..a279d7e2bee 100644 --- a/vertx-core/src/main/java/io/vertx/core/http/impl/http1/Http1ServerConnection.java +++ b/vertx-core/src/main/java/io/vertx/core/http/impl/http1/Http1ServerConnection.java @@ -38,7 +38,7 @@ import io.vertx.core.internal.http.QueryParamDecoder; import io.vertx.core.net.NetSocket; import io.vertx.core.net.ServerSSLOptions; -import io.vertx.core.net.impl.MessageWrite; +import io.vertx.core.net.impl.WritePromise; import io.vertx.core.net.impl.tcp.NetSocketImpl; import io.vertx.core.internal.tls.SslContextManager; import io.vertx.core.net.impl.VertxHandler; @@ -263,20 +263,18 @@ private void onEnd() { } } - void write(VertxHttpObject msg, Promise promise) { - writeToChannel(new MessageWrite() { + Future write(VertxHttpObject msg) { + WritePromise write = new WritePromise(context) { @Override public void write() { - Http1ServerConnection.this.unsafeWrite(msg, false, promise); + Http1ServerConnection.this.unsafeWrite(msg, false, this); if (msg.isEnded()) { responseComplete(); } } - @Override - public void cancel(Throwable cause) { - promise.fail(cause); - } - }); + }; + writeToChannel(write); + return write; } void responseComplete() { @@ -487,25 +485,25 @@ protected void handleWriteQueueDrained() { } } - void write100Continue(Promise promise) { - write(new VertxFullHttpResponse( + Future write100Continue() { + return write(new VertxFullHttpResponse( false, HTTP_1_1, CONTINUE, Unpooled.buffer(0), DefaultHttpHeadersFactory.headersFactory().newHeaders(), DefaultHttpHeadersFactory.trailersFactory().newHeaders(), - false), promise); + false)); } - void write103EarlyHints(HttpHeaders headers, Promise promise) { - write(new VertxFullHttpResponse(false, + Future write103EarlyHints(HttpHeaders headers) { + return write(new VertxFullHttpResponse(false, HTTP_1_1, HttpResponseStatus.EARLY_HINTS, Unpooled.buffer(0), headers, EmptyHttpHeaders.INSTANCE, - false), promise); + false)); } protected void handleClosed() { diff --git a/vertx-core/src/main/java/io/vertx/core/http/impl/http1/Http1ServerResponse.java b/vertx-core/src/main/java/io/vertx/core/http/impl/http1/Http1ServerResponse.java index 9cc8b7e6103..221a86b2a3b 100644 --- a/vertx-core/src/main/java/io/vertx/core/http/impl/http1/Http1ServerResponse.java +++ b/vertx-core/src/main/java/io/vertx/core/http/impl/http1/Http1ServerResponse.java @@ -312,7 +312,7 @@ public HttpServerResponse endHandler(@Nullable Handler handler) { @Override public Future writeHead() { checkThread(); - PromiseInternal promise = context.promise(); + Future f; synchronized (conn) { if (headWritten) { throw new IllegalStateException(); @@ -323,44 +323,35 @@ public Future writeHead() { VertxHttpObject msg; prepareHeaders(-1); msg = new VertxHttpResponse(head, version, status, headers); - conn.write(msg, promise); + f = conn.write(msg); } - return promise.future(); + return f; } @Override public Future write(Buffer chunk) { - PromiseInternal promise = context.promise(); - write(((BufferInternal)chunk).getByteBuf(), promise); - return promise.future(); + return write(((BufferInternal)chunk).getByteBuf()); } @Override public Future write(String chunk, String enc) { - PromiseInternal promise = context.promise(); - write(BufferInternal.buffer(chunk, enc).getByteBuf(), promise); - return promise.future(); + return write(BufferInternal.buffer(chunk, enc).getByteBuf()); } @Override public Future write(String chunk) { - PromiseInternal promise = context.promise(); - write(BufferInternal.buffer(chunk).getByteBuf(), promise); - return promise.future(); + return write(BufferInternal.buffer(chunk).getByteBuf()); } @Override public Future writeContinue() { checkThread(); - Promise promise = context.promise(); - conn.write100Continue(promise); - return promise.future(); + return conn.write100Continue(); } @Override public Future writeEarlyHints(MultiMap headers) { checkThread(); - PromiseInternal promise = context.promise(); Http1xHeaders headersMultiMap; if (headers instanceof Http1xHeaders) { headersMultiMap = (Http1xHeaders) headers; @@ -371,8 +362,7 @@ public Future writeEarlyHints(MultiMap headers) { synchronized (conn) { checkHeadWritten(); } - conn.write103EarlyHints(headersMultiMap, promise); - return promise.future(); + return conn.write103EarlyHints(headersMultiMap); } @Override @@ -387,12 +377,6 @@ public Future end(String chunk, String enc) { @Override public Future end(Buffer chunk) { - PromiseInternal promise = context.promise(); - end(chunk, promise); - return promise.future(); - } - - private void end(Buffer chunk, PromiseInternal listener) { checkThread(); synchronized (conn) { if (written) { @@ -410,7 +394,7 @@ private void end(Buffer chunk, PromiseInternal listener) { } else { msg = new VertxLastHttpContent(data, trailingHeaders); } - conn.write(msg, listener); + Future result = conn.write(msg); if (bodyEndHandler != null) { bodyEndHandler.handle(null); } @@ -420,6 +404,7 @@ private void end(Buffer chunk, PromiseInternal listener) { if (!keepAlive) { closed = true; // ????? } + return result; } } @@ -514,7 +499,7 @@ private Future sendFileInternal(long offset, long length, long size, Rando prepareHeaders(actualLength); bytesWritten = actualLength; written = true; - conn.write(new VertxAssembledHttpResponse(head, version, status, headers), null); + conn.write(new VertxAssembledHttpResponse(head, version, status, headers)); FileChannel toSend = fileChannel == null ? file.getChannel() : fileChannel; ChannelFuture channelFuture = conn.sendFile(toSend, actualOffset, actualLength); PromiseInternal promise = context.promise(); @@ -541,7 +526,8 @@ private Future sendFileInternal(long offset, long length, long size, Rando } // write an empty last content to let the http encoder know the response is complete - conn.write(new VertxLastHttpContent(Unpooled.buffer(0), DefaultHttpHeadersFactory.trailersFactory().newHeaders()), promise); + Future f = conn.write(new VertxLastHttpContent(Unpooled.buffer(0), DefaultHttpHeadersFactory.trailersFactory().newHeaders())); + f.onComplete(promise); } else { promise.fail(future.cause()); } @@ -718,7 +704,7 @@ private void reportResponseBegin() { } } - private Http1ServerResponse write(ByteBuf chunk, PromiseInternal promise) { + private Future write(ByteBuf chunk) { checkThread(); synchronized (conn) { if (written) { @@ -736,8 +722,7 @@ private Http1ServerResponse write(ByteBuf chunk, PromiseInternal promise) } else { msg = new VertxHttpContent(chunk); } - conn.write(msg, promise); - return this; + return conn.write(msg); } } @@ -758,8 +743,7 @@ Future netSocket(HttpMethod requestMethod, MultiMap requestHeaders) { } status = requestMethod == HttpMethod.CONNECT ? HttpResponseStatus.OK : HttpResponseStatus.SWITCHING_PROTOCOLS; prepareHeaders(-1); - PromiseInternal upgradePromise = context.promise(); - conn.write(new VertxAssembledHttpResponse(head, version, status, headers), upgradePromise); + conn.write(new VertxAssembledHttpResponse(head, version, status, headers)); written = true; Promise promise = context.promise(); netSocket = promise.future(); diff --git a/vertx-core/src/main/java/io/vertx/core/http/impl/http2/DefaultHttp2ClientStream.java b/vertx-core/src/main/java/io/vertx/core/http/impl/http2/DefaultHttp2ClientStream.java index 6c84b39f2d0..6b7c1e27af2 100644 --- a/vertx-core/src/main/java/io/vertx/core/http/impl/http2/DefaultHttp2ClientStream.java +++ b/vertx-core/src/main/java/io/vertx/core/http/impl/http2/DefaultHttp2ClientStream.java @@ -35,7 +35,7 @@ import io.vertx.core.internal.PromiseInternal; import io.vertx.core.internal.buffer.BufferInternal; import io.vertx.core.net.HostAndPort; -import io.vertx.core.net.impl.MessageWrite; +import io.vertx.core.net.impl.WritePromise; import io.vertx.core.spi.metrics.ClientMetrics; import io.vertx.core.spi.metrics.TransportMetrics; import io.vertx.core.spi.tracing.VertxTracer; @@ -117,7 +117,7 @@ public Future writeHead(HttpRequestHead request, boolean chunked, Buffer b priority(priority); scheme = request.scheme; authority = request.authority; - write(new HeadersWrite(request, buf != null ? ((BufferInternal)buf).getByteBuf() : null, end, promise)); + write(new HeadersWrite(context, request, buf != null ? ((BufferInternal)buf).getByteBuf() : null, end, promise)); return promise.future(); } @@ -128,17 +128,18 @@ public HttpClientStream setWriteQueueMaxSize(int maxSize) { void writeHeaders(HttpRequestHead request, ByteBuf buf, boolean end, StreamPriority priority, Promise promise) { priority(priority); - write(new HeadersWrite(request, buf, end, promise)); + write(new HeadersWrite(context, request, buf, end, promise)); } - private class HeadersWrite implements MessageWrite { + private class HeadersWrite extends WritePromise { private final HttpRequestHead request; private final ByteBuf buf; private final boolean end; private final Promise promise; - public HeadersWrite(HttpRequestHead request, ByteBuf buf, boolean end, Promise promise) { + public HeadersWrite(ContextInternal context, HttpRequestHead request, ByteBuf buf, boolean end, Promise promise) { + super(context); this.request = request; this.buf = buf; this.end = end; diff --git a/vertx-core/src/main/java/io/vertx/core/http/impl/http2/DefaultHttp2ServerStream.java b/vertx-core/src/main/java/io/vertx/core/http/impl/http2/DefaultHttp2ServerStream.java index 54ed55c170c..3664481e6d8 100644 --- a/vertx-core/src/main/java/io/vertx/core/http/impl/http2/DefaultHttp2ServerStream.java +++ b/vertx-core/src/main/java/io/vertx/core/http/impl/http2/DefaultHttp2ServerStream.java @@ -164,16 +164,14 @@ void handleHeader(HttpHeaders map) { } public final Future writeHead(HttpResponseHead head, Buffer chunk, boolean end) { - Promise promise = context.promise(); HttpResponseHeaders headers = (HttpResponseHeaders)head.headers(); headers.status(head.statusCode); if (chunk != null) { - writeHeaders(headers, false, false, null); - writeData(((BufferInternal)chunk).getByteBuf(), end, promise); + writeHeaders(headers, false, false); + return writeData(((BufferInternal)chunk).getByteBuf(), end); } else { - writeHeaders(headers, end, true, promise); + return writeHeaders(headers, end, true); } - return promise.future(); } public void routed(String route) { diff --git a/vertx-core/src/main/java/io/vertx/core/http/impl/http2/DefaultHttp2Stream.java b/vertx-core/src/main/java/io/vertx/core/http/impl/http2/DefaultHttp2Stream.java index 2ce861bb6bf..046a292bb85 100644 --- a/vertx-core/src/main/java/io/vertx/core/http/impl/http2/DefaultHttp2Stream.java +++ b/vertx-core/src/main/java/io/vertx/core/http/impl/http2/DefaultHttp2Stream.java @@ -32,7 +32,7 @@ import io.vertx.core.internal.buffer.BufferInternal; import io.vertx.core.internal.concurrent.InboundMessageQueue; import io.vertx.core.internal.concurrent.OutboundMessageQueue; -import io.vertx.core.net.impl.MessageWrite; +import io.vertx.core.net.impl.WritePromise; /** * @author Julien Viet @@ -41,7 +41,7 @@ abstract class DefaultHttp2Stream> implements Ht private static final HttpHeaders EMPTY = new HttpHeaders(EmptyHttp2Headers.INSTANCE); - private final OutboundMessageQueue outboundQueue; + private final OutboundMessageQueue outboundQueue; private final InboundMessageQueue inboundQueue; private final Http2Connection connection; protected final VertxInternal vertx; @@ -100,7 +100,7 @@ protected void handleMessage(Object item) { this.outboundQueue = new OutboundMessageQueue<>(connection.context().executor()) { // TODO implement stop drain to optimize flushes ? @Override - public boolean test(MessageWrite msg) { + public boolean test(WritePromise msg) { if (DefaultHttp2Stream.this.writable) { msg.write(); return true; @@ -109,7 +109,7 @@ public boolean test(MessageWrite msg) { } } @Override - protected void handleDispose(MessageWrite messageWrite) { + protected void handleDispose(WritePromise messageWrite) { Throwable cause = failure; if (cause == null) { cause = HttpUtils.STREAM_CLOSED_EXCEPTION; @@ -268,8 +268,10 @@ public final boolean isWritable() { return outboundQueue.isWritable(); } - public final void write(MessageWrite write) { - outboundQueue.write(write); + public final Future write(WritePromise write) { + boolean writable = outboundQueue.write(write); + write.enqueued(writable); + return write; } public final S pause() { @@ -295,32 +297,28 @@ public final Future writeFrame(int type, int flags, Buffer payload) { } public final Future writeHeaders(MultiMap headers, boolean end) { - Promise promise = context.promise(); - writeHeaders((HttpHeaders) headers, end, true, promise); - return promise.future(); + return writeHeaders((HttpHeaders) headers, end, true); } - void writeHeaders(HttpHeaders headers, boolean end, boolean checkFlush, Promise promise) { + Future writeHeaders(HttpHeaders headers, boolean end, boolean checkFlush) { + WritePromise write = new WritePromise(context) { + @Override + public void write() { + writeHeaders0(headers, end, checkFlush, this); + } + }; if (first_) { first_ = false; EventLoop eventLoop = connection.context().nettyEventLoop(); if (eventLoop.inEventLoop()) { - writeHeaders0(headers, end, checkFlush, promise); + write.write(); } else { - eventLoop.execute(() -> writeHeaders0(headers, end, checkFlush, promise)); + eventLoop.execute(write::write); } } else { - outboundQueue.write(new MessageWrite() { - @Override - public void write() { - writeHeaders0(headers, end, checkFlush, promise); - } - @Override - public void cancel(Throwable cause) { - promise.fail(cause); - } - }); + write(write); } + return write; } void writeHeaders0(HttpHeaders headers, boolean end, boolean checkFlush, Promise promise) { @@ -356,7 +354,7 @@ void writeHeaders0(HttpHeaders headers, boolean end, boolean checkFlush, Promise public final void sendFile(ChunkedInput file, Promise promise) { bytesWritten += file.length(); - outboundQueue.write(new MessageWrite() { + write(new WritePromise(context) { @Override public void write() { sendFile0(file, promise); @@ -373,20 +371,14 @@ private void sendFile0(ChunkedInput file, Promise promise) { } public final Future writeChunk(Buffer chunk, boolean end) { - Promise promise = context.promise(); - writeData(chunk == null ? null : ((BufferInternal)chunk).getByteBuf(), end, promise); - return promise.future(); + return writeData(chunk == null ? null : ((BufferInternal)chunk).getByteBuf(), end); } - public final void writeData(ByteBuf chunk, boolean end, Promise promise) { - write(new MessageWrite() { + public final Future writeData(ByteBuf chunk, boolean end) { + return write(new WritePromise(context) { @Override public void write() { - writeData0(chunk == null ? Unpooled.EMPTY_BUFFER : chunk, end, promise); - } - @Override - public void cancel(Throwable cause) { - promise.fail(cause); + writeData0(chunk == null ? Unpooled.EMPTY_BUFFER : chunk, end, this); } }); } diff --git a/vertx-core/src/main/java/io/vertx/core/impl/future/PromiseImpl.java b/vertx-core/src/main/java/io/vertx/core/impl/future/PromiseImpl.java index dc3e3e3d5db..ec1d2b0638d 100644 --- a/vertx-core/src/main/java/io/vertx/core/impl/future/PromiseImpl.java +++ b/vertx-core/src/main/java/io/vertx/core/impl/future/PromiseImpl.java @@ -20,7 +20,7 @@ * * @author Julien Viet */ -public final class PromiseImpl extends FutureImpl implements PromiseInternal { +public class PromiseImpl extends FutureImpl implements PromiseInternal { /** * Create a promise that hasn't completed yet diff --git a/vertx-core/src/main/java/io/vertx/core/internal/streams/WriteResult.java b/vertx-core/src/main/java/io/vertx/core/internal/streams/WriteResult.java new file mode 100644 index 00000000000..702b9194b86 --- /dev/null +++ b/vertx-core/src/main/java/io/vertx/core/internal/streams/WriteResult.java @@ -0,0 +1,17 @@ +package io.vertx.core.internal.streams; + +import io.vertx.core.Future; + +/** + * Result of a write operation. + * + * @author Julien Viet + */ +public interface WriteResult extends Future { + + /** + * @return whether the stream was writable after the write operation was submitted + */ + boolean isWritable(); + +} diff --git a/vertx-core/src/main/java/io/vertx/core/net/impl/MessageWrite.java b/vertx-core/src/main/java/io/vertx/core/net/impl/MessageWrite.java index afa4953f3f3..59fc2222079 100644 --- a/vertx-core/src/main/java/io/vertx/core/net/impl/MessageWrite.java +++ b/vertx-core/src/main/java/io/vertx/core/net/impl/MessageWrite.java @@ -21,6 +21,9 @@ public interface MessageWrite { */ void write(); + default void enqueued(boolean writable) { + } + /** * Cancel the write operation. * diff --git a/vertx-core/src/main/java/io/vertx/core/net/impl/StreamChannelBase.java b/vertx-core/src/main/java/io/vertx/core/net/impl/StreamChannelBase.java index 788524231af..c0d66e50d48 100644 --- a/vertx-core/src/main/java/io/vertx/core/net/impl/StreamChannelBase.java +++ b/vertx-core/src/main/java/io/vertx/core/net/impl/StreamChannelBase.java @@ -22,7 +22,6 @@ import io.vertx.codegen.annotations.Nullable; import io.vertx.core.Future; import io.vertx.core.Handler; -import io.vertx.core.Promise; import io.vertx.core.ThreadingModel; import io.vertx.core.buffer.Buffer; import io.vertx.core.impl.EventLoopExecutor; @@ -90,9 +89,14 @@ protected void handleMessage(Object msg) { @Override public Future writeMessage(Object message) { - Promise promise = context.promise(); - writeToChannel(message, promise); - return promise.future(); + WritePromise messageWrite = new WritePromise(context) { + @Override + public void write() { + unsafeWrite(message, false, this); + } + }; + writeToChannel(messageWrite); + return messageWrite; } @Override diff --git a/vertx-core/src/main/java/io/vertx/core/net/impl/VertxConnection.java b/vertx-core/src/main/java/io/vertx/core/net/impl/VertxConnection.java index cb946abddbf..299da415928 100644 --- a/vertx-core/src/main/java/io/vertx/core/net/impl/VertxConnection.java +++ b/vertx-core/src/main/java/io/vertx/core/net/impl/VertxConnection.java @@ -481,7 +481,9 @@ public void cancel(Throwable cause) { // Write to channel boolean return for now is not used so avoids reading a volatile public final boolean writeToChannel(MessageWrite msg) { - return outboundMessageQueue.write(msg); + boolean writable = outboundMessageQueue.write(msg); + msg.enqueued(writable); + return writable; } /** diff --git a/vertx-core/src/main/java/io/vertx/core/net/impl/WritePromise.java b/vertx-core/src/main/java/io/vertx/core/net/impl/WritePromise.java new file mode 100644 index 00000000000..f2ce5080bb0 --- /dev/null +++ b/vertx-core/src/main/java/io/vertx/core/net/impl/WritePromise.java @@ -0,0 +1,43 @@ +/* + * Copyright (c) 2011-2026 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0 + * which is available at https://www.apache.org/licenses/LICENSE-2.0. + * + * SPDX-License-Identifier: EPL-2.0 OR Apache-2.0 + */ +package io.vertx.core.net.impl; + +import io.vertx.core.impl.future.PromiseImpl; +import io.vertx.core.internal.ContextInternal; +import io.vertx.core.internal.streams.WriteResult; + +/** + * A write promise. + * + * @author Julien Viet + */ +public abstract class WritePromise extends PromiseImpl implements MessageWrite, WriteResult { + + private boolean writable; + + public WritePromise(ContextInternal context) { + super(context); + } + + @Override + public void enqueued(boolean writable) { + this.writable = writable; + } + + @Override + public void cancel(Throwable cause) { + fail(cause); + } + + public boolean isWritable() { + return writable; + } +} diff --git a/vertx-core/src/main/java/io/vertx/core/net/impl/quic/QuicStreamImpl.java b/vertx-core/src/main/java/io/vertx/core/net/impl/quic/QuicStreamImpl.java index 18a607b1c68..774eed12a12 100644 --- a/vertx-core/src/main/java/io/vertx/core/net/impl/quic/QuicStreamImpl.java +++ b/vertx-core/src/main/java/io/vertx/core/net/impl/quic/QuicStreamImpl.java @@ -21,7 +21,7 @@ import io.vertx.core.internal.ContextInternal; import io.vertx.core.internal.PromiseInternal; import io.vertx.core.internal.net.QuicStreamInternal; -import io.vertx.core.net.impl.MessageWrite; +import io.vertx.core.net.impl.WritePromise; import io.vertx.core.net.impl.StreamChannelBase; import io.vertx.core.net.QuicConnection; import io.vertx.core.net.QuicStream; @@ -110,8 +110,7 @@ protected long sizeof(Object msg) { @Override public Future end() { - PromiseInternal promise = context.promise(); - writeToChannel(new MessageWrite() { + WritePromise write = new WritePromise(context) { @Override public void write() { ChannelFuture shutdownPromise; @@ -122,14 +121,11 @@ public void write() { } else { shutdownPromise = channel.shutdownOutput(); } - shutdownPromise.addListener(promise); - } - @Override - public void cancel(Throwable cause) { - promise.fail(cause); + shutdownPromise.addListener(this); } - }); - return promise.future(); + }; + writeToChannel(write); + return write; } @Override diff --git a/vertx-core/src/test/java/io/vertx/tests/http/HttpTest.java b/vertx-core/src/test/java/io/vertx/tests/http/HttpTest.java index d0a4741c991..5a98f5c304d 100644 --- a/vertx-core/src/test/java/io/vertx/tests/http/HttpTest.java +++ b/vertx-core/src/test/java/io/vertx/tests/http/HttpTest.java @@ -32,8 +32,10 @@ import io.vertx.core.internal.ContextInternal; import io.vertx.core.internal.http.HttpClientInternal; import io.vertx.core.internal.net.endpoint.EndpointResolverInternal; +import io.vertx.core.internal.streams.WriteResult; import io.vertx.core.json.JsonObject; import io.vertx.core.net.*; +import io.vertx.core.net.impl.MessageWrite; import io.vertx.core.streams.ReadStream; import io.vertx.test.core.*; import io.vertx.test.fakedns.DnsRecord; @@ -6662,4 +6664,77 @@ public Buffer testNoAuthority(boolean force) throws Exception { }) .await(); } + + @Test + public void testServerMessageWrite() throws Exception { + AtomicReference> ref1 = new AtomicReference<>(); + AtomicReference> ref2 = new AtomicReference<>(); + server.requestHandler(request -> { + HttpServerResponse response = request + .response() + .setChunked(true); + ref1.set(response.write("chunk")); + ref2.set(response.end("last")); + }); + startServer(); + client.request(new RequestOptions(requestOptions).setPort(server.actualPort())) + .compose(req -> req + .send() + .expecting(HttpResponseExpectation.SC_OK) + .compose(HttpClientResponse::body)); + TestUtils.assertWaitUntil(() -> ref1.get() != null); + TestUtils.assertWaitUntil(() -> ref2.get() != null); + assertTrue(ref1.get() instanceof MessageWrite); + assertTrue(ref2.get() instanceof MessageWrite); + } + + private void stableUnwritable(Buffer chunk, HttpServerResponse response, WriteResult prev, Consumer> done) { + WriteResult last = null; + while (true) { + if (response.writeQueueFull()) { + WriteResult w = last; + if (w == null) { + // Done => use prev + done.accept(prev); + } else { + // Give some time to fill the window + vertx.setTimer(10, v -> { + stableUnwritable(chunk, response, w, done); + }); + } + break; + } else { + last = (WriteResult)response.write(chunk); + } + } + } + + @Test + public void testServerMessageWriteWritability(Checkpoint checkpoint1, Checkpoint checkpoint2, Checkpoint checkpoint3) throws Exception { + server.requestHandler(request -> { + HttpServerResponse response = request.response(); + response.setChunked(true); + vertx.runOnContext(v1 -> { + Buffer chunk = Buffer.buffer(TestUtils.randomAlphaString(1024)); + stableUnwritable(chunk, response, null, last -> { + assertEquals(last.isWritable(), !response.writeQueueFull()); + last.onComplete(v -> { + checkpoint2.succeed(); + }); + response.drainHandler(v2 -> { + checkpoint3.succeed(); + }); + checkpoint1.succeed(); + }); + }); + }); + startServer(); + HttpClientResponse response = client.request(new RequestOptions(requestOptions).setPort(server.actualPort())) + .compose(req -> req + .send() + .expecting(HttpResponseExpectation.SC_OK) + .andThen(onSuccess(HttpClientResponse::pause))).await(); + checkpoint1.awaitSuccess(); + response.resume(); + } }