diff --git a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/AbstractH2StreamMultiplexer.java b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/AbstractH2StreamMultiplexer.java index dadbf2799..81057d804 100644 --- a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/AbstractH2StreamMultiplexer.java +++ b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/AbstractH2StreamMultiplexer.java @@ -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 it = streams.iterator(); it.hasNext(); ) { + it.next().resetIfCancelled(); + } } final int pendingOutputRequests = outputRequests.get(); boolean outputPending = false; @@ -1314,6 +1319,7 @@ private void consumeSettingsFrame(final ByteBuffer payload) throws IOException { private void produceOutput() throws HttpException, IOException { for (final Iterator it = streams.iterator(); it.hasNext(); ) { final H2Stream stream = it.next(); + stream.resetIfCancelled(); if (!stream.isLocalClosed() && !stream.isReserved() && stream.getOutputWindow().get() > 0) { stream.produceOutput(); } diff --git a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/H2Stream.java b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/H2Stream.java index f946a3665..88081f695 100644 --- a/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/H2Stream.java +++ b/httpcore5-h2/src/main/java/org/apache/hc/core5/http2/impl/nio/H2Stream.java @@ -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); } diff --git a/httpcore5-h2/src/test/java/org/apache/hc/core5/http2/impl/nio/TestAbstractH2StreamMultiplexer.java b/httpcore5-h2/src/test/java/org/apache/hc/core5/http2/impl/nio/TestAbstractH2StreamMultiplexer.java index 4b72cfec3..411e1ac0b 100644 --- a/httpcore5-h2/src/test/java/org/apache/hc/core5/http2/impl/nio/TestAbstractH2StreamMultiplexer.java +++ b/httpcore5-h2/src/test/java/org/apache/hc/core5/http2/impl/nio/TestAbstractH2StreamMultiplexer.java @@ -1183,6 +1183,118 @@ void testStreamIdleTimeoutTriggersH2StreamTimeoutException() throws Exception { } + @Test + void testAbortAfterLocalEndStreamSendsRstStreamWithoutInboundFrames() throws Exception { + final List 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
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 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 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
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 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();