From 6c062149fc23066043668d9ad43e9c4b0fef2a59 Mon Sep 17 00:00:00 2001 From: dragonfsky Date: Sun, 23 Aug 2026 11:14:06 +0800 Subject: [PATCH 1/3] Issue #14403 - Defer FCGI application dispatch until the parser returns Signed-off-by: Dongliang Xie --- .../server/internal/HttpStreamOverFCGI.java | 13 +- .../server/internal/ServerFCGIConnection.java | 27 +- .../internal/ServerFCGIConnectionTest.java | 345 ++++++++++++++++++ 3 files changed, 371 insertions(+), 14 deletions(-) create mode 100644 jetty-core/jetty-fcgi/jetty-fcgi-server/src/test/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnectionTest.java diff --git a/jetty-core/jetty-fcgi/jetty-fcgi-server/src/main/java/org/eclipse/jetty/fcgi/server/internal/HttpStreamOverFCGI.java b/jetty-core/jetty-fcgi/jetty-fcgi-server/src/main/java/org/eclipse/jetty/fcgi/server/internal/HttpStreamOverFCGI.java index 575fdd44d912..d0a7213fa16b 100644 --- a/jetty-core/jetty-fcgi/jetty-fcgi-server/src/main/java/org/eclipse/jetty/fcgi/server/internal/HttpStreamOverFCGI.java +++ b/jetty-core/jetty-fcgi/jetty-fcgi-server/src/main/java/org/eclipse/jetty/fcgi/server/internal/HttpStreamOverFCGI.java @@ -103,17 +103,15 @@ else if (FCGI.Headers.HTTPS.equalsIgnoreCase(name)) processField(field); } - public void onHeaders() + public Runnable onHeaders() { String pathQuery = URIUtil.addPathQuery(_path, _query); HttpScheme scheme = StringUtil.isEmpty(_secure) ? HttpScheme.HTTP : HttpScheme.HTTPS; MetaData.Request request = new MetaData.Request(_connection.getBeginNanoTime(), _method, scheme.asString(), hostPort, pathQuery, HttpVersion.fromString(_version), _headers, -1); Runnable task = _httpChannel.onRequest(request); _allHeaders.forEach(field -> _httpChannel.getRequest().setAttribute(field.getName(), field.getValue())); - // TODO: here we just execute the task. - // However, we should really return all the way back to onFillable() - // and feed the Runnable to an ExecutionStrategy. - execute(task); + // Return the task to dispatch it after ServerParser.parse() has returned. + return task; } private void processField(HttpField field) @@ -358,11 +356,6 @@ public boolean onIdleTimeout(TimeoutException timeout) return !handlingRequest; } - private void execute(Runnable task) - { - _connection.getConnector().getExecutor().execute(task); - } - private class DemandCallback implements Callback { @Override diff --git a/jetty-core/jetty-fcgi/jetty-fcgi-server/src/main/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnection.java b/jetty-core/jetty-fcgi/jetty-fcgi-server/src/main/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnection.java index b0460dd7c629..9a9d3972f246 100644 --- a/jetty-core/jetty-fcgi/jetty-fcgi-server/src/main/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnection.java +++ b/jetty-core/jetty-fcgi/jetty-fcgi-server/src/main/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnection.java @@ -15,6 +15,7 @@ import java.nio.ByteBuffer; import java.util.Set; +import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.TimeoutException; import org.eclipse.jetty.fcgi.FCGI; @@ -59,6 +60,7 @@ public class ServerFCGIConnection extends AbstractMetaDataConnection implements private boolean useOutputDirectByteBuffers; private RetainableByteBuffer inputBuffer; private HttpStreamOverFCGI stream; + private Runnable onRequest; public ServerFCGIConnection(Connector connector, EndPoint endPoint, HttpConfiguration configuration, boolean sendStatus200) { @@ -191,7 +193,7 @@ public void onFillable() { if (stream == null && inputBuffer.isEmpty()) releaseInputBuffer(); - return; + break; } } else if (read == 0) @@ -207,6 +209,24 @@ else if (read == 0) return; } } + + // Dispatch only after the parser has returned and input buffer bookkeeping is complete. + Runnable task = onRequest; + onRequest = null; + if (task != null) + { + try + { + getExecutor().execute(task); + } + catch (RejectedExecutionException x) + { + HttpStreamOverFCGI stream = this.stream; + Runnable failureTask = stream.getHttpChannel().onFailure(x); + this.stream = null; + ThreadPool.executeImmediately(getExecutor(), failureTask); + } + } } catch (Exception x) { @@ -359,9 +379,8 @@ public boolean onHeaders(int request) LOG.debug("Request {} headers on {}", request, stream); if (stream != null) { - stream.onHeaders(); - // We have dispatched to the application, - // so we must stop the fill & parse loop. + onRequest = stream.onHeaders(); + // Return to onFillable() before dispatching to the application. return true; } return false; diff --git a/jetty-core/jetty-fcgi/jetty-fcgi-server/src/test/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnectionTest.java b/jetty-core/jetty-fcgi/jetty-fcgi-server/src/test/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnectionTest.java new file mode 100644 index 000000000000..6014891240bb --- /dev/null +++ b/jetty-core/jetty-fcgi/jetty-fcgi-server/src/test/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnectionTest.java @@ -0,0 +1,345 @@ +// +// ======================================================================== +// Copyright (c) 1995 Mort Bay Consulting Pty Ltd and others. +// +// This program and the accompanying materials are made available under the +// terms of the Eclipse Public License v. 2.0 which is available at +// https://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 org.eclipse.jetty.fcgi.server.internal; + +import java.lang.reflect.Field; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executor; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; + +import org.eclipse.jetty.fcgi.FCGI; +import org.eclipse.jetty.fcgi.generator.ClientGenerator; +import org.eclipse.jetty.fcgi.parser.ClientParser; +import org.eclipse.jetty.fcgi.server.ServerFCGIConnectionFactory; +import org.eclipse.jetty.http.HttpFields; +import org.eclipse.jetty.http.HttpStatus; +import org.eclipse.jetty.http.HttpVersion; +import org.eclipse.jetty.io.ArrayByteBufferPool; +import org.eclipse.jetty.io.ByteArrayEndPoint; +import org.eclipse.jetty.io.ByteBufferPool; +import org.eclipse.jetty.io.RetainableByteBuffer; +import org.eclipse.jetty.logging.StacklessLogging; +import org.eclipse.jetty.server.Handler; +import org.eclipse.jetty.server.HttpConfiguration; +import org.eclipse.jetty.server.Request; +import org.eclipse.jetty.server.Response; +import org.eclipse.jetty.server.Server; +import org.eclipse.jetty.server.ServerConnector; +import org.eclipse.jetty.util.BufferUtil; +import org.eclipse.jetty.util.Callback; +import org.junit.jupiter.api.Test; + +import static org.awaitility.Awaitility.await; +import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.is; +import static org.hamcrest.Matchers.nullValue; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +public class ServerFCGIConnectionTest +{ + /** + * The application completes while its first input-buffer release is parked, + * reproducing the input-buffer ownership handoff of issue #14403. + */ + @Test + public void testApplicationDispatchedAfterParsing() throws Exception + { + AtomicInteger handlersInvoked = new AtomicInteger(); + AtomicInteger completedTasks = new AtomicInteger(); + AtomicReference taskFailure = new AtomicReference<>(); + AtomicBoolean firstAcquire = new AtomicBoolean(true); + AtomicBoolean releasedOnce = new AtomicBoolean(); + CountDownLatch inputReleaseEntered = new CountDownLatch(1); + CountDownLatch inputReleaseProceed = new CountDownLatch(1); + + ArrayByteBufferPool.Tracking trackingPool = new ArrayByteBufferPool.Tracking(); + ByteBufferPool bufferPool = new ByteBufferPool.Wrapper(trackingPool) + { + @Override + public RetainableByteBuffer.Mutable acquire(int size, boolean direct) + { + RetainableByteBuffer.Mutable buffer = super.acquire(size, direct); + if (firstAcquire.getAndSet(false)) + { + return new RetainableByteBuffer.Mutable.Wrapper(buffer) + { + @Override + public boolean release() + { + if (releasedOnce.compareAndSet(false, true)) + { + inputReleaseEntered.countDown(); + awaitLatch(inputReleaseProceed, "input buffer release was not released"); + } + return super.release(); + } + }; + } + return buffer; + } + }; + + ExecutorService handlerExecutor = Executors.newSingleThreadExecutor(); + Executor executor = task -> + { + handlerExecutor.execute(() -> + { + try + { + task.run(); + } + catch (Throwable x) + { + taskFailure.set(x); + } + finally + { + completedTasks.incrementAndGet(); + } + }); + awaitLatch(inputReleaseEntered, "application task did not reach the input buffer release"); + }; + + ByteArrayEndPoint endPoint = new ByteArrayEndPoint(new byte[0], 64 * 1024); + Server server = newServer(handlersInvoked); + ServerFCGIConnection connection = newConnection(server, executor, bufferPool, endPoint); + + endPoint.addInput(generateRequest(1)); + + AtomicReference fillableFailure = new AtomicReference<>(); + Thread fillThread = new Thread(() -> run(connection::onFillable, fillableFailure), "fcgi-io"); + fillThread.start(); + try + { + fillThread.join(TimeUnit.SECONDS.toMillis(10)); + assertFalse(fillThread.isAlive(), "onFillable() did not complete"); + + inputReleaseProceed.countDown(); + await().atMost(10, TimeUnit.SECONDS).until(() -> completedTasks.get() == 1); + assertEquals(1, handlersInvoked.get()); + assertThat(fillableFailure.get(), nullValue()); + + Responses responses = new Responses(); + parseResponse(endPoint.takeOutput(), responses); + responses.assertResponses(200, 1); + + endPoint.addInput(generateRequest(2)); + await().atMost(10, TimeUnit.SECONDS).until(() -> completedTasks.get() == 2); + assertThat(taskFailure.get(), nullValue()); + parseResponse(endPoint.takeOutput(), responses); + responses.assertResponses(200, 2); + assertEquals(2, handlersInvoked.get()); + + assertThat("Server Leaks: " + trackingPool.dumpLeaks(), trackingPool.getLeaks().size(), is(0)); + } + finally + { + inputReleaseProceed.countDown(); + fillThread.join(TimeUnit.SECONDS.toMillis(10)); + handlerExecutor.shutdownNow(); + assertTrue(handlerExecutor.awaitTermination(10, TimeUnit.SECONDS)); + } + } + + @Test + public void testApplicationDispatchedAfterParsingWithInlineExecutor() throws Exception + { + AtomicInteger handlersInvoked = new AtomicInteger(); + ArrayByteBufferPool.Tracking trackingPool = new ArrayByteBufferPool.Tracking(); + ByteArrayEndPoint endPoint = new ByteArrayEndPoint(new byte[0], 64 * 1024); + Server server = newServer(handlersInvoked); + ServerFCGIConnection connection = newConnection(server, Runnable::run, trackingPool, endPoint); + + endPoint.addInput(generateRequest(1)); + + connection.onFillable(); + assertEquals(1, handlersInvoked.get()); + + Responses responses = new Responses(); + parseResponse(endPoint.takeOutput(), responses); + responses.assertResponses(HttpStatus.OK_200, 1); + assertThat("Server Leaks: " + trackingPool.dumpLeaks(), trackingPool.getLeaks().size(), is(0)); + } + + @Test + public void testApplicationDispatchRejectedBeforeRequestEnd() throws Exception + { + AtomicInteger handlersInvoked = new AtomicInteger(); + ArrayByteBufferPool.Tracking trackingPool = new ArrayByteBufferPool.Tracking(); + Executor executor = task -> + { + throw new RejectedExecutionException("test"); + }; + ByteArrayEndPoint endPoint = new ByteArrayEndPoint(new byte[0], 64 * 1024); + Server server = newServer(handlersInvoked); + ServerFCGIConnection connection = newConnection(server, executor, trackingPool, endPoint); + + endPoint.addInput(generateRequestHeaders(1)); + + ByteBuffer output; + try (StacklessLogging ignored = new StacklessLogging(Response.class)) + { + connection.onFillable(); + output = endPoint.waitForOutput(10, TimeUnit.SECONDS); + } + assertEquals(0, handlersInvoked.get()); + assertTrue(output != null && output.hasRemaining(), "no failure response written"); + // The failure must terminate the connection-side request even if its content has not arrived. + Field streamField = ServerFCGIConnection.class.getDeclaredField("stream"); + streamField.setAccessible(true); + assertThat(streamField.get(connection), nullValue()); + + Responses responses = new Responses(); + parseResponse(output, responses); + responses.assertResponses(HttpStatus.INTERNAL_SERVER_ERROR_500, 1); + await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> + assertThat("Server Leaks: " + trackingPool.dumpLeaks(), trackingPool.getLeaks().size(), is(0))); + } + + private static Server newServer(AtomicInteger handlersInvoked) + { + Server server = new Server(); + server.setHandler(new Handler.Abstract() + { + @Override + public boolean handle(Request request, Response response, Callback callback) + { + handlersInvoked.incrementAndGet(); + // Do not read the request content, just complete the callback + // as soon as the request arrives, like the application in + // issue #14403. + callback.succeeded(); + return true; + } + }); + return server; + } + + private static ServerFCGIConnection newConnection(Server server, Executor executor, ByteBufferPool bufferPool, ByteArrayEndPoint endPoint) + { + ServerFCGIConnectionFactory connectionFactory = new ServerFCGIConnectionFactory(new HttpConfiguration()); + ServerConnector connector = new ServerConnector(server, executor, null, bufferPool, 0, 1, connectionFactory); + return new ServerFCGIConnection(connector, endPoint, new HttpConfiguration(), false); + } + + private static class Responses implements ClientParser.Listener + { + private final List statuses = new ArrayList<>(); + private final AtomicInteger headers = new AtomicInteger(); + private final AtomicInteger ends = new AtomicInteger(); + + @Override + public void onBegin(int request, int code, String reason) + { + statuses.add(code); + } + + @Override + public boolean onHeaders(int request) + { + headers.incrementAndGet(); + return false; + } + + @Override + public boolean onEnd(int request) + { + ends.incrementAndGet(); + return true; + } + + void assertResponses(int status, int count) + { + assertEquals(count, statuses.size()); + for (int i = 0; i < count; ++i) + assertEquals(status, statuses.get(i)); + assertEquals(count, headers.get()); + assertEquals(count, ends.get()); + } + } + + private static void parseResponse(ByteBuffer output, ClientParser.Listener listener) + { + assertTrue(output.hasRemaining(), "no response written"); + new ClientParser(listener).parse(output); + } + + private static ByteBuffer generateRequest(int id) + { + return generateRequest(id, true); + } + + private static ByteBuffer generateRequestHeaders(int id) + { + return generateRequest(id, false); + } + + private static ByteBuffer generateRequest(int id, boolean complete) + { + ClientGenerator generator = new ClientGenerator(ByteBufferPool.NON_POOLING); + ByteBufferPool.Accumulator accumulator = new ByteBufferPool.Accumulator(); + HttpFields.Mutable params = HttpFields.build() + .put(FCGI.Headers.REQUEST_METHOD, "GET") + .put(FCGI.Headers.DOCUMENT_URI, "/") + .put(FCGI.Headers.QUERY_STRING, "") + .put(FCGI.Headers.SERVER_PROTOCOL, HttpVersion.HTTP_1_1.asString()); + generator.generateRequestHeaders(accumulator, id, params); + if (complete) + generator.generateRequestContent(accumulator, id, BufferUtil.EMPTY_BUFFER, true); + List buffers = accumulator.getByteBuffers(); + int capacity = (int)accumulator.getTotalLength(); + ByteBuffer request = ByteBuffer.allocate(capacity); + buffers.forEach(request::put); + accumulator.release(); + BufferUtil.flipToFlush(request, 0); + return request; + } + + private static void run(Runnable action, AtomicReference failure) + { + try + { + action.run(); + } + catch (Throwable x) + { + failure.set(x); + } + } + + private static void awaitLatch(CountDownLatch latch, String message) + { + try + { + if (!latch.await(10, TimeUnit.SECONDS)) + throw new IllegalStateException(message); + } + catch (InterruptedException x) + { + Thread.currentThread().interrupt(); + throw new RuntimeException(x); + } + } +} From 63f2e774858b8aef0b984436e8a6fb3badfeaf0e Mon Sep 17 00:00:00 2001 From: dragonfsky Date: Sun, 23 Aug 2026 20:05:46 +0800 Subject: [PATCH 2/3] Rerun CI Signed-off-by: Dongliang Xie From 402fb008ef8a57f53916f4beecbff9b79e400189 Mon Sep 17 00:00:00 2001 From: dragonfsky Date: Mon, 24 Aug 2026 18:55:42 +0800 Subject: [PATCH 3/3] Issue #14403 - Simplify FCGI dispatch fix Signed-off-by: Dongliang Xie --- .../server/internal/ServerFCGIConnection.java | 25 +- .../internal/ServerFCGIConnectionTest.java | 284 +----------------- 2 files changed, 17 insertions(+), 292 deletions(-) diff --git a/jetty-core/jetty-fcgi/jetty-fcgi-server/src/main/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnection.java b/jetty-core/jetty-fcgi/jetty-fcgi-server/src/main/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnection.java index 9a9d3972f246..faf3bb44ac21 100644 --- a/jetty-core/jetty-fcgi/jetty-fcgi-server/src/main/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnection.java +++ b/jetty-core/jetty-fcgi/jetty-fcgi-server/src/main/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnection.java @@ -15,7 +15,6 @@ import java.nio.ByteBuffer; import java.util.Set; -import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.TimeoutException; import org.eclipse.jetty.fcgi.FCGI; @@ -193,7 +192,11 @@ public void onFillable() { if (stream == null && inputBuffer.isEmpty()) releaseInputBuffer(); - break; + Runnable task = onRequest; + onRequest = null; + if (task != null) + getExecutor().execute(task); + return; } } else if (read == 0) @@ -209,24 +212,6 @@ else if (read == 0) return; } } - - // Dispatch only after the parser has returned and input buffer bookkeeping is complete. - Runnable task = onRequest; - onRequest = null; - if (task != null) - { - try - { - getExecutor().execute(task); - } - catch (RejectedExecutionException x) - { - HttpStreamOverFCGI stream = this.stream; - Runnable failureTask = stream.getHttpChannel().onFailure(x); - this.stream = null; - ThreadPool.executeImmediately(getExecutor(), failureTask); - } - } } catch (Exception x) { diff --git a/jetty-core/jetty-fcgi/jetty-fcgi-server/src/test/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnectionTest.java b/jetty-core/jetty-fcgi/jetty-fcgi-server/src/test/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnectionTest.java index 6014891240bb..64bd0530b009 100644 --- a/jetty-core/jetty-fcgi/jetty-fcgi-server/src/test/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnectionTest.java +++ b/jetty-core/jetty-fcgi/jetty-fcgi-server/src/test/java/org/eclipse/jetty/fcgi/server/internal/ServerFCGIConnectionTest.java @@ -13,32 +13,18 @@ package org.eclipse.jetty.fcgi.server.internal; -import java.lang.reflect.Field; import java.nio.ByteBuffer; -import java.util.ArrayList; import java.util.List; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.Executor; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.RejectedExecutionException; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; -import java.util.concurrent.atomic.AtomicReference; import org.eclipse.jetty.fcgi.FCGI; import org.eclipse.jetty.fcgi.generator.ClientGenerator; -import org.eclipse.jetty.fcgi.parser.ClientParser; import org.eclipse.jetty.fcgi.server.ServerFCGIConnectionFactory; import org.eclipse.jetty.http.HttpFields; -import org.eclipse.jetty.http.HttpStatus; import org.eclipse.jetty.http.HttpVersion; import org.eclipse.jetty.io.ArrayByteBufferPool; import org.eclipse.jetty.io.ByteArrayEndPoint; import org.eclipse.jetty.io.ByteBufferPool; -import org.eclipse.jetty.io.RetainableByteBuffer; -import org.eclipse.jetty.logging.StacklessLogging; import org.eclipse.jetty.server.Handler; import org.eclipse.jetty.server.HttpConfiguration; import org.eclipse.jetty.server.Request; @@ -49,177 +35,14 @@ import org.eclipse.jetty.util.Callback; import org.junit.jupiter.api.Test; -import static org.awaitility.Awaitility.await; -import static org.hamcrest.MatcherAssert.assertThat; -import static org.hamcrest.Matchers.is; -import static org.hamcrest.Matchers.nullValue; import static org.junit.jupiter.api.Assertions.assertEquals; -import static org.junit.jupiter.api.Assertions.assertFalse; -import static org.junit.jupiter.api.Assertions.assertTrue; public class ServerFCGIConnectionTest { - /** - * The application completes while its first input-buffer release is parked, - * reproducing the input-buffer ownership handoff of issue #14403. - */ @Test - public void testApplicationDispatchedAfterParsing() throws Exception + public void testApplicationDispatchedAfterParsing() { AtomicInteger handlersInvoked = new AtomicInteger(); - AtomicInteger completedTasks = new AtomicInteger(); - AtomicReference taskFailure = new AtomicReference<>(); - AtomicBoolean firstAcquire = new AtomicBoolean(true); - AtomicBoolean releasedOnce = new AtomicBoolean(); - CountDownLatch inputReleaseEntered = new CountDownLatch(1); - CountDownLatch inputReleaseProceed = new CountDownLatch(1); - - ArrayByteBufferPool.Tracking trackingPool = new ArrayByteBufferPool.Tracking(); - ByteBufferPool bufferPool = new ByteBufferPool.Wrapper(trackingPool) - { - @Override - public RetainableByteBuffer.Mutable acquire(int size, boolean direct) - { - RetainableByteBuffer.Mutable buffer = super.acquire(size, direct); - if (firstAcquire.getAndSet(false)) - { - return new RetainableByteBuffer.Mutable.Wrapper(buffer) - { - @Override - public boolean release() - { - if (releasedOnce.compareAndSet(false, true)) - { - inputReleaseEntered.countDown(); - awaitLatch(inputReleaseProceed, "input buffer release was not released"); - } - return super.release(); - } - }; - } - return buffer; - } - }; - - ExecutorService handlerExecutor = Executors.newSingleThreadExecutor(); - Executor executor = task -> - { - handlerExecutor.execute(() -> - { - try - { - task.run(); - } - catch (Throwable x) - { - taskFailure.set(x); - } - finally - { - completedTasks.incrementAndGet(); - } - }); - awaitLatch(inputReleaseEntered, "application task did not reach the input buffer release"); - }; - - ByteArrayEndPoint endPoint = new ByteArrayEndPoint(new byte[0], 64 * 1024); - Server server = newServer(handlersInvoked); - ServerFCGIConnection connection = newConnection(server, executor, bufferPool, endPoint); - - endPoint.addInput(generateRequest(1)); - - AtomicReference fillableFailure = new AtomicReference<>(); - Thread fillThread = new Thread(() -> run(connection::onFillable, fillableFailure), "fcgi-io"); - fillThread.start(); - try - { - fillThread.join(TimeUnit.SECONDS.toMillis(10)); - assertFalse(fillThread.isAlive(), "onFillable() did not complete"); - - inputReleaseProceed.countDown(); - await().atMost(10, TimeUnit.SECONDS).until(() -> completedTasks.get() == 1); - assertEquals(1, handlersInvoked.get()); - assertThat(fillableFailure.get(), nullValue()); - - Responses responses = new Responses(); - parseResponse(endPoint.takeOutput(), responses); - responses.assertResponses(200, 1); - - endPoint.addInput(generateRequest(2)); - await().atMost(10, TimeUnit.SECONDS).until(() -> completedTasks.get() == 2); - assertThat(taskFailure.get(), nullValue()); - parseResponse(endPoint.takeOutput(), responses); - responses.assertResponses(200, 2); - assertEquals(2, handlersInvoked.get()); - - assertThat("Server Leaks: " + trackingPool.dumpLeaks(), trackingPool.getLeaks().size(), is(0)); - } - finally - { - inputReleaseProceed.countDown(); - fillThread.join(TimeUnit.SECONDS.toMillis(10)); - handlerExecutor.shutdownNow(); - assertTrue(handlerExecutor.awaitTermination(10, TimeUnit.SECONDS)); - } - } - - @Test - public void testApplicationDispatchedAfterParsingWithInlineExecutor() throws Exception - { - AtomicInteger handlersInvoked = new AtomicInteger(); - ArrayByteBufferPool.Tracking trackingPool = new ArrayByteBufferPool.Tracking(); - ByteArrayEndPoint endPoint = new ByteArrayEndPoint(new byte[0], 64 * 1024); - Server server = newServer(handlersInvoked); - ServerFCGIConnection connection = newConnection(server, Runnable::run, trackingPool, endPoint); - - endPoint.addInput(generateRequest(1)); - - connection.onFillable(); - assertEquals(1, handlersInvoked.get()); - - Responses responses = new Responses(); - parseResponse(endPoint.takeOutput(), responses); - responses.assertResponses(HttpStatus.OK_200, 1); - assertThat("Server Leaks: " + trackingPool.dumpLeaks(), trackingPool.getLeaks().size(), is(0)); - } - - @Test - public void testApplicationDispatchRejectedBeforeRequestEnd() throws Exception - { - AtomicInteger handlersInvoked = new AtomicInteger(); - ArrayByteBufferPool.Tracking trackingPool = new ArrayByteBufferPool.Tracking(); - Executor executor = task -> - { - throw new RejectedExecutionException("test"); - }; - ByteArrayEndPoint endPoint = new ByteArrayEndPoint(new byte[0], 64 * 1024); - Server server = newServer(handlersInvoked); - ServerFCGIConnection connection = newConnection(server, executor, trackingPool, endPoint); - - endPoint.addInput(generateRequestHeaders(1)); - - ByteBuffer output; - try (StacklessLogging ignored = new StacklessLogging(Response.class)) - { - connection.onFillable(); - output = endPoint.waitForOutput(10, TimeUnit.SECONDS); - } - assertEquals(0, handlersInvoked.get()); - assertTrue(output != null && output.hasRemaining(), "no failure response written"); - // The failure must terminate the connection-side request even if its content has not arrived. - Field streamField = ServerFCGIConnection.class.getDeclaredField("stream"); - streamField.setAccessible(true); - assertThat(streamField.get(connection), nullValue()); - - Responses responses = new Responses(); - parseResponse(output, responses); - responses.assertResponses(HttpStatus.INTERNAL_SERVER_ERROR_500, 1); - await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> - assertThat("Server Leaks: " + trackingPool.dumpLeaks(), trackingPool.getLeaks().size(), is(0))); - } - - private static Server newServer(AtomicInteger handlersInvoked) - { Server server = new Server(); server.setHandler(new Handler.Abstract() { @@ -227,77 +50,17 @@ private static Server newServer(AtomicInteger handlersInvoked) public boolean handle(Request request, Response response, Callback callback) { handlersInvoked.incrementAndGet(); - // Do not read the request content, just complete the callback - // as soon as the request arrives, like the application in - // issue #14403. callback.succeeded(); return true; } }); - return server; - } - private static ServerFCGIConnection newConnection(Server server, Executor executor, ByteBufferPool bufferPool, ByteArrayEndPoint endPoint) - { + ArrayByteBufferPool.Tracking bufferPool = new ArrayByteBufferPool.Tracking(); ServerFCGIConnectionFactory connectionFactory = new ServerFCGIConnectionFactory(new HttpConfiguration()); - ServerConnector connector = new ServerConnector(server, executor, null, bufferPool, 0, 1, connectionFactory); - return new ServerFCGIConnection(connector, endPoint, new HttpConfiguration(), false); - } - - private static class Responses implements ClientParser.Listener - { - private final List statuses = new ArrayList<>(); - private final AtomicInteger headers = new AtomicInteger(); - private final AtomicInteger ends = new AtomicInteger(); - - @Override - public void onBegin(int request, int code, String reason) - { - statuses.add(code); - } - - @Override - public boolean onHeaders(int request) - { - headers.incrementAndGet(); - return false; - } - - @Override - public boolean onEnd(int request) - { - ends.incrementAndGet(); - return true; - } - - void assertResponses(int status, int count) - { - assertEquals(count, statuses.size()); - for (int i = 0; i < count; ++i) - assertEquals(status, statuses.get(i)); - assertEquals(count, headers.get()); - assertEquals(count, ends.get()); - } - } - - private static void parseResponse(ByteBuffer output, ClientParser.Listener listener) - { - assertTrue(output.hasRemaining(), "no response written"); - new ClientParser(listener).parse(output); - } - - private static ByteBuffer generateRequest(int id) - { - return generateRequest(id, true); - } - - private static ByteBuffer generateRequestHeaders(int id) - { - return generateRequest(id, false); - } + ServerConnector connector = new ServerConnector(server, Runnable::run, null, bufferPool, 0, 1, connectionFactory); + ByteArrayEndPoint endPoint = new ByteArrayEndPoint(new byte[0], 64 * 1024); + ServerFCGIConnection connection = new ServerFCGIConnection(connector, endPoint, new HttpConfiguration(), false); - private static ByteBuffer generateRequest(int id, boolean complete) - { ClientGenerator generator = new ClientGenerator(ByteBufferPool.NON_POOLING); ByteBufferPool.Accumulator accumulator = new ByteBufferPool.Accumulator(); HttpFields.Mutable params = HttpFields.build() @@ -305,41 +68,18 @@ private static ByteBuffer generateRequest(int id, boolean complete) .put(FCGI.Headers.DOCUMENT_URI, "/") .put(FCGI.Headers.QUERY_STRING, "") .put(FCGI.Headers.SERVER_PROTOCOL, HttpVersion.HTTP_1_1.asString()); - generator.generateRequestHeaders(accumulator, id, params); - if (complete) - generator.generateRequestContent(accumulator, id, BufferUtil.EMPTY_BUFFER, true); + generator.generateRequestHeaders(accumulator, 1, params); + generator.generateRequestContent(accumulator, 1, BufferUtil.EMPTY_BUFFER, true); List buffers = accumulator.getByteBuffers(); - int capacity = (int)accumulator.getTotalLength(); - ByteBuffer request = ByteBuffer.allocate(capacity); + ByteBuffer request = ByteBuffer.allocate((int)accumulator.getTotalLength()); buffers.forEach(request::put); accumulator.release(); BufferUtil.flipToFlush(request, 0); - return request; - } - private static void run(Runnable action, AtomicReference failure) - { - try - { - action.run(); - } - catch (Throwable x) - { - failure.set(x); - } - } + endPoint.addInput(request); + connection.onFillable(); - private static void awaitLatch(CountDownLatch latch, String message) - { - try - { - if (!latch.await(10, TimeUnit.SECONDS)) - throw new IllegalStateException(message); - } - catch (InterruptedException x) - { - Thread.currentThread().interrupt(); - throw new RuntimeException(x); - } + assertEquals(1, handlersInvoked.get()); + assertEquals(0, bufferPool.getLeaks().size(), bufferPool.dumpLeaks()); } }