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 @@ -526,6 +526,11 @@ public final void onOutput() throws HttpException, IOException {

if (connOutputWindow.get() > 0 && remoteSettingState == SettingsHandshake.ACKED) {
produceOutput();
} else {
// RST_STREAM is not subject to flow control
for (final Iterator<H2Stream> it = streams.iterator(); it.hasNext(); ) {
it.next().resetIfCancelled();
}
}
final int pendingOutputRequests = outputRequests.get();
boolean outputPending = false;
Expand Down Expand Up @@ -1314,6 +1319,7 @@ private void consumeSettingsFrame(final ByteBuffer payload) throws IOException {
private void produceOutput() throws HttpException, IOException {
for (final Iterator<H2Stream> it = streams.iterator(); it.hasNext(); ) {
final H2Stream stream = it.next();
stream.resetIfCancelled();
if (!stream.isLocalClosed() && !stream.isReserved() && stream.getOutputWindow().get() > 0) {
stream.produceOutput();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -266,6 +266,16 @@ void localResetCancelled() throws IOException {
localReset(new H2StreamResetException(H2Error.CANCEL, "Cancelled"));
}

/**
* Resets a cancelled stream that has already sent END_STREAM. Such a stream
* produces no more output, so it would otherwise wait for a peer frame.
*/
void resetIfCancelled() throws IOException {
if (cancelled.get() && channel.isLocalClosed() && !isClosed()) {
localResetCancelled();
}
}

void handle(final HttpException ex) throws IOException, HttpException {
handler.handle(ex, remoteClosed);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1183,6 +1183,118 @@ void testStreamIdleTimeoutTriggersH2StreamTimeoutException() throws Exception {

}

@Test
void testAbortAfterLocalEndStreamSendsRstStreamWithoutInboundFrames() throws Exception {
final List<byte[]> writes = new ArrayList<>();
Mockito.when(protocolIOSession.write(ArgumentMatchers.any(ByteBuffer.class)))
.thenAnswer(inv -> {
final ByteBuffer b = inv.getArgument(0, ByteBuffer.class);
final byte[] copy = new byte[b.remaining()];
b.get(copy);
writes.add(copy);
return copy.length;
});
Mockito.doNothing().when(protocolIOSession).setEvent(ArgumentMatchers.anyInt());
Mockito.doNothing().when(protocolIOSession).clearEvent(ArgumentMatchers.anyInt());

final H2Config h2Config = H2Config.custom().build();
final H2StreamMultiplexerImpl mux = new H2StreamMultiplexerImpl(
protocolIOSession, FRAME_FACTORY, StreamIdGenerator.ODD,
httpProcessor, CharCodingConfig.DEFAULT, h2Config, h2StreamListener, () -> streamHandler);

mux.onConnect();
final WritableByteChannelMock writable = new WritableByteChannelMock(256);
final FrameOutputBuffer fob = new FrameOutputBuffer(16 * 1024);
fob.write(new RawFrame(FrameType.SETTINGS.getValue(), 0, 0, null), writable);
mux.onInput(ByteBuffer.wrap(writable.toByteArray()));
writes.clear();

// A request without a body: HEADERS carry END_STREAM, so the stream is half-closed (local)
final H2StreamChannel channel = mux.createChannel(1);
final List<Header> requestHeaders = Arrays.asList(
new BasicHeader(":method", "GET"),
new BasicHeader(":scheme", "https"),
new BasicHeader(":path", "/"),
new BasicHeader(":authority", "example.test"));
final H2Stream stream = mux.createStream(channel, new PriorityHeaderSender(channel, requestHeaders, true));
mux.onOutput();
Assertions.assertTrue(stream.isLocalClosed());
writes.clear();

// The request gets cancelled while the peer stays silent
stream.abort();
mux.onOutput();

final List<FrameStub> frames = parseFrames(concat(writes));
final FrameStub rst = frames.stream()
.filter(f -> f.type == FrameType.RST_STREAM.getValue() && f.streamId == 1)
.findFirst()
.orElse(null);
Assertions.assertNotNull(rst, "RST_STREAM not emitted for the cancelled stream");
Assertions.assertEquals(H2Error.CANCEL.getCode(), ByteBuffer.wrap(rst.payload).getInt());
Assertions.assertTrue(channel.isLocalReset());
}

@Test
void testAbortAfterLocalEndStreamSendsRstStreamWithNoConnectionWindow() throws Exception {
final List<byte[]> writes = new ArrayList<>();
Mockito.when(protocolIOSession.write(ArgumentMatchers.any(ByteBuffer.class)))
.thenAnswer(inv -> {
final ByteBuffer b = inv.getArgument(0, ByteBuffer.class);
final byte[] copy = new byte[b.remaining()];
b.get(copy);
writes.add(copy);
return copy.length;
});
Mockito.doNothing().when(protocolIOSession).setEvent(ArgumentMatchers.anyInt());
Mockito.doNothing().when(protocolIOSession).clearEvent(ArgumentMatchers.anyInt());

final H2Config h2Config = H2Config.custom().build();
final H2StreamMultiplexerImpl mux = new H2StreamMultiplexerImpl(
protocolIOSession, FRAME_FACTORY, StreamIdGenerator.ODD,
httpProcessor, CharCodingConfig.DEFAULT, h2Config, h2StreamListener, () -> streamHandler);

mux.onConnect();
final WritableByteChannelMock writable = new WritableByteChannelMock(256);
final FrameOutputBuffer fob = new FrameOutputBuffer(16 * 1024);
fob.write(new RawFrame(FrameType.SETTINGS.getValue(), 0, 0, null), writable);
mux.onInput(ByteBuffer.wrap(writable.toByteArray()));

final List<Header> requestHeaders = Arrays.asList(
new BasicHeader(":method", "GET"),
new BasicHeader(":scheme", "https"),
new BasicHeader(":path", "/"),
new BasicHeader(":authority", "example.test"));
// Stream 1: a request without a body, half-closed (local) once its HEADERS are sent
final H2StreamChannel channel = mux.createChannel(1);
final H2Stream stream = mux.createStream(channel, new PriorityHeaderSender(channel, requestHeaders, true));
// Stream 3: a request with a body that uses up the connection output window
final H2StreamChannel uploadChannel = mux.createChannel(3);
mux.createStream(uploadChannel, new PriorityHeaderSender(uploadChannel, requestHeaders, false));
mux.onOutput();
Assertions.assertTrue(stream.isLocalClosed());

final ByteBuffer body = ByteBuffer.allocate(h2Config.getInitialWindowSize());
while (body.hasRemaining()) {
Assertions.assertTrue(uploadChannel.write(body) > 0, "connection window used up too early");
}
Assertions.assertEquals(0, uploadChannel.write(ByteBuffer.allocate(1)), "connection window must be 0");
writes.clear();

// The request gets cancelled while the peer stays silent and grants no window
stream.abort();
mux.onOutput();

final List<FrameStub> frames = parseFrames(concat(writes));
final FrameStub rst = frames.stream()
.filter(f -> f.type == FrameType.RST_STREAM.getValue() && f.streamId == 1)
.findFirst()
.orElse(null);
Assertions.assertNotNull(rst, "RST_STREAM not emitted while the connection window is 0");
Assertions.assertEquals(H2Error.CANCEL.getCode(), ByteBuffer.wrap(rst.payload).getInt());
Assertions.assertTrue(channel.isLocalReset());
}

@Test
void testResetIfExpiredResetsStreamPastDeadline() throws Exception {
final H2Config h2Config = H2Config.custom().build();
Expand Down
Loading