diff --git a/spring-boot-project/spring-boot/src/main/java/org/springframework/boot/web/embedded/jetty/JettyReactiveWebServerFactory.java b/spring-boot-project/spring-boot/src/main/java/org/springframework/boot/web/embedded/jetty/JettyReactiveWebServerFactory.java index b386bbaab0..531667de61 100644 --- a/spring-boot-project/spring-boot/src/main/java/org/springframework/boot/web/embedded/jetty/JettyReactiveWebServerFactory.java +++ b/spring-boot-project/spring-boot/src/main/java/org/springframework/boot/web/embedded/jetty/JettyReactiveWebServerFactory.java @@ -172,6 +172,7 @@ public class JettyReactiveWebServerFactory extends AbstractReactiveWebServerFact InetSocketAddress address = new InetSocketAddress(getAddress(), port); Server server = new Server(getThreadPool()); server.addConnector(createConnector(address, server)); + server.setStopTimeout(0); ServletHolder servletHolder = new ServletHolder(servlet); servletHolder.setAsyncSupported(true); ServletContextHandler contextHandler = new ServletContextHandler(server, "/", false, false); diff --git a/spring-boot-project/spring-boot/src/main/java/org/springframework/boot/web/embedded/jetty/JettyServletWebServerFactory.java b/spring-boot-project/spring-boot/src/main/java/org/springframework/boot/web/embedded/jetty/JettyServletWebServerFactory.java index 2385cf663d..d53ef01a11 100644 --- a/spring-boot-project/spring-boot/src/main/java/org/springframework/boot/web/embedded/jetty/JettyServletWebServerFactory.java +++ b/spring-boot-project/spring-boot/src/main/java/org/springframework/boot/web/embedded/jetty/JettyServletWebServerFactory.java @@ -173,6 +173,7 @@ public class JettyServletWebServerFactory extends AbstractServletWebServerFactor private Server createServer(InetSocketAddress address) { Server server = new Server(getThreadPool()); server.setConnectors(new Connector[] { createConnector(address, server) }); + server.setStopTimeout(0); return server; } diff --git a/spring-boot-project/spring-boot/src/test/java/org/springframework/boot/web/reactive/server/AbstractReactiveWebServerFactoryTests.java b/spring-boot-project/spring-boot/src/test/java/org/springframework/boot/web/reactive/server/AbstractReactiveWebServerFactoryTests.java index 5f95b61c2b..f8ece0a60d 100644 --- a/spring-boot-project/spring-boot/src/test/java/org/springframework/boot/web/reactive/server/AbstractReactiveWebServerFactoryTests.java +++ b/spring-boot-project/spring-boot/src/test/java/org/springframework/boot/web/reactive/server/AbstractReactiveWebServerFactoryTests.java @@ -431,6 +431,30 @@ public abstract class AbstractReactiveWebServerFactoryTests { blockingHandler.completeOne(); } + @Test + void whenARequestIsActiveAfterGracefulShutdownEndsThenStopWillComplete() throws InterruptedException { + AbstractReactiveWebServerFactory factory = getFactory(); + factory.setShutdown(Shutdown.GRACEFUL); + BlockingHandler blockingHandler = new BlockingHandler(); + this.webServer = factory.getWebServer(blockingHandler); + this.webServer.start(); + Mono> request = getWebClient(this.webServer.getPort()).build().get().retrieve() + .toBodilessEntity(); + AtomicReference> responseReference = new AtomicReference<>(); + CountDownLatch responseLatch = new CountDownLatch(1); + request.subscribe((response) -> { + responseReference.set(response); + responseLatch.countDown(); + }); + blockingHandler.awaitQueue(); + AtomicReference result = new AtomicReference<>(); + this.webServer.shutDownGracefully(result::set); + this.webServer.stop(); + Awaitility.await().atMost(Duration.ofSeconds(30)) + .until(() -> GracefulShutdownResult.REQUESTS_ACTIVE == result.get()); + blockingHandler.completeOne(); + } + @Test void whenARequestIsActiveThenStopWillComplete() throws InterruptedException, BrokenBarrierException { AbstractReactiveWebServerFactory factory = getFactory(); diff --git a/spring-boot-project/spring-boot/src/test/java/org/springframework/boot/web/servlet/server/AbstractServletWebServerFactoryTests.java b/spring-boot-project/spring-boot/src/test/java/org/springframework/boot/web/servlet/server/AbstractServletWebServerFactoryTests.java index 9e1096101a..74bb1ad2e0 100644 --- a/spring-boot-project/spring-boot/src/test/java/org/springframework/boot/web/servlet/server/AbstractServletWebServerFactoryTests.java +++ b/spring-boot-project/spring-boot/src/test/java/org/springframework/boot/web/servlet/server/AbstractServletWebServerFactoryTests.java @@ -1159,6 +1159,31 @@ public abstract class AbstractServletWebServerFactoryTests { assertThat(getResponse("http://localhost:" + this.webServer.getPort() + "/hello")).isEqualTo("Hello World"); } + @Test + void whenARequestIsActiveAfterGracefulShutdownEndsThenStopWillComplete() + throws InterruptedException, BrokenBarrierException { + AbstractServletWebServerFactory factory = getFactory(); + factory.setShutdown(Shutdown.GRACEFUL); + BlockingServlet blockingServlet = new BlockingServlet(); + this.webServer = factory + .getWebServer((context) -> context.addServlet("blockingServlet", blockingServlet).addMapping("/")); + this.webServer.start(); + int port = this.webServer.getPort(); + initiateGetRequest(port, "/"); + blockingServlet.awaitQueue(); + AtomicReference result = new AtomicReference<>(); + this.webServer.shutDownGracefully(result::set); + this.webServer.stop(); + Awaitility.await().atMost(Duration.ofSeconds(30)) + .until(() -> GracefulShutdownResult.REQUESTS_ACTIVE == result.get()); + try { + blockingServlet.admitOne(); + } + catch (RuntimeException ex) { + + } + } + protected Future initiateGetRequest(int port, String path) { return initiateGetRequest(HttpClients.createMinimal(), port, path); }