This is an automated email from the ASF dual-hosted git repository. FrankChen021 pushed a commit to branch codex/response-server-headers in repository https://gitbox.apache.org/repos/asf/druid.git
commit cb191b827fb6d0cd5b9a756c71e18a9916122ae9 Author: Frank Chen <[email protected]> AuthorDate: Sat Sep 5 22:46:54 2026 +0800 fix(server): identify locally generated proxy errors --- .../server/ResponseIdentityHeaderTest.java | 44 +++++++++++++++++++++ .../server/AsyncManagementForwardingServlet.java | 13 +++++++ .../druid/server/http/OverlordProxyServlet.java | 13 +++++++ .../initialization/jetty/JettyServerModule.java | 11 ++++++ .../jetty/ResponseIdentityHeaderHandler.java | 45 ++++++++++++++++++++++ .../jetty/ResponseIdentityHeaderHandlerTest.java | 27 +++++++++++++ .../druid/server/AsyncQueryForwardingServlet.java | 14 +++++++ 7 files changed, 167 insertions(+) diff --git a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/server/ResponseIdentityHeaderTest.java b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/server/ResponseIdentityHeaderTest.java index fa956d10d85..7a279978cfb 100644 --- a/embedded-tests/src/test/java/org/apache/druid/testing/embedded/server/ResponseIdentityHeaderTest.java +++ b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/server/ResponseIdentityHeaderTest.java @@ -34,9 +34,11 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; import java.net.URI; +import java.net.Socket; import java.net.http.HttpClient; import java.net.http.HttpRequest; import java.net.http.HttpResponse; +import java.nio.charset.StandardCharsets; import java.time.Duration; import java.util.List; @@ -135,6 +137,40 @@ public class ResponseIdentityHeaderTest extends EmbeddedClusterTestBase assertResponseIdentity(client.send(request, HttpResponse.BodyHandlers.ofString()), router, 405); } + @Test + @Timeout(30) + public void testJettyRequestParsingErrorUsesRouterIdentity() throws Exception + { + final URI routerUri = URI.create(getServerUrl(router)); + final String response; + try (Socket socket = new Socket(routerUri.getHost(), routerUri.getPort())) { + socket.setSoTimeout(10_000); + socket.getOutputStream().write( + "GET /status/health HTTP/1.1\r\nHost: localhost\r\nInvalid Header: value\r\n\r\n" + .getBytes(StandardCharsets.US_ASCII) + ); + socket.getOutputStream().flush(); + response = new String(socket.getInputStream().readAllBytes(), StandardCharsets.US_ASCII); + } + + Assertions.assertTrue(response.startsWith("HTTP/1.1 400"), response); + assertRawHeader( + response, + ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER, + router.bindings().selfNode().getHostAndPortToUse() + ); + assertRawHeader( + response, + ResponseIdentityHeaderHandler.RESPONSE_SERVICE_HEADER, + router.bindings().selfNode().getServiceName() + ); + assertRawHeader( + response, + ResponseIdentityHeaderHandler.RESPONSE_VERSION_HEADER, + router.bindings().selfNode().getVersion() + ); + } + private HttpResponse<String> sendGet(final String url) throws Exception { final HttpRequest request = HttpRequest.newBuilder(URI.create(url)) @@ -212,4 +248,12 @@ public class ResponseIdentityHeaderTest extends EmbeddedClusterTestBase response.headers().allValues(ResponseIdentityHeaderHandler.RESPONSE_VERSION_HEADER) ); } + + private static void assertRawHeader(final String response, final String name, final String value) + { + Assertions.assertTrue( + response.lines().anyMatch(line -> line.equalsIgnoreCase(name + ": " + value)), + response + ); + } } diff --git a/server/src/main/java/org/apache/druid/server/AsyncManagementForwardingServlet.java b/server/src/main/java/org/apache/druid/server/AsyncManagementForwardingServlet.java index 757eaf510d7..042b0232fff 100644 --- a/server/src/main/java/org/apache/druid/server/AsyncManagementForwardingServlet.java +++ b/server/src/main/java/org/apache/druid/server/AsyncManagementForwardingServlet.java @@ -202,11 +202,24 @@ public class AsyncManagementForwardingServlet extends AsyncProxyServlet Response serverResponse ) { + ResponseIdentityHeaderHandler.rememberLocalIdentity(clientRequest, proxyResponse); ResponseIdentityHeaderHandler.clearRouterIdentity(proxyResponse); StandardResponseHeaderFilterHolder.deduplicateHeadersInProxyServlet(proxyResponse, serverResponse); super.onServerResponseHeaders(clientRequest, proxyResponse, serverResponse); } + @Override + protected void onProxyResponseFailure( + final HttpServletRequest clientRequest, + final HttpServletResponse proxyResponse, + final Response serverResponse, + final Throwable failure + ) + { + ResponseIdentityHeaderHandler.restoreLocalIdentity(clientRequest, proxyResponse); + super.onProxyResponseFailure(clientRequest, proxyResponse, serverResponse, failure); + } + @Override protected HttpField filterServerResponseHeader( final HttpServletRequest clientRequest, diff --git a/server/src/main/java/org/apache/druid/server/http/OverlordProxyServlet.java b/server/src/main/java/org/apache/druid/server/http/OverlordProxyServlet.java index fabdab23d69..5de54650dfa 100644 --- a/server/src/main/java/org/apache/druid/server/http/OverlordProxyServlet.java +++ b/server/src/main/java/org/apache/druid/server/http/OverlordProxyServlet.java @@ -98,11 +98,24 @@ public class OverlordProxyServlet extends ProxyServlet final Response serverResponse ) { + ResponseIdentityHeaderHandler.rememberLocalIdentity(clientRequest, proxyResponse); ResponseIdentityHeaderHandler.clearRouterIdentity(proxyResponse); StandardResponseHeaderFilterHolder.deduplicateHeadersInProxyServlet(proxyResponse, serverResponse); super.onServerResponseHeaders(clientRequest, proxyResponse, serverResponse); } + @Override + protected void onProxyResponseFailure( + final HttpServletRequest clientRequest, + final HttpServletResponse proxyResponse, + final Response serverResponse, + final Throwable failure + ) + { + ResponseIdentityHeaderHandler.restoreLocalIdentity(clientRequest, proxyResponse); + super.onProxyResponseFailure(clientRequest, proxyResponse, serverResponse, failure); + } + @Override protected HttpField filterServerResponseHeader( final HttpServletRequest clientRequest, diff --git a/server/src/main/java/org/apache/druid/server/initialization/jetty/JettyServerModule.java b/server/src/main/java/org/apache/druid/server/initialization/jetty/JettyServerModule.java index b6246ac9a1a..f534a41f604 100644 --- a/server/src/main/java/org/apache/druid/server/initialization/jetty/JettyServerModule.java +++ b/server/src/main/java/org/apache/druid/server/initialization/jetty/JettyServerModule.java @@ -485,6 +485,17 @@ public class JettyServerModule extends JerseyServletModule }); } + if (config.isEnableResponseIdentityHeaders()) { + // Request parsing failures do not enter the server's handler chain, so add the identity at the error handler too. + final Request.Handler errorHandler = server.getErrorHandler() == null + ? new ErrorHandler() + : server.getErrorHandler(); + server.setErrorHandler((request, response, callback) -> { + ResponseIdentityHeaderHandler.addIdentityHeaders(response, node); + return errorHandler.handle(request, response, callback); + }); + } + server.setRequestLog(new JettyRequestLog()); return server; diff --git a/server/src/main/java/org/apache/druid/server/initialization/jetty/ResponseIdentityHeaderHandler.java b/server/src/main/java/org/apache/druid/server/initialization/jetty/ResponseIdentityHeaderHandler.java index d9860e60998..a6e8a722014 100644 --- a/server/src/main/java/org/apache/druid/server/initialization/jetty/ResponseIdentityHeaderHandler.java +++ b/server/src/main/java/org/apache/druid/server/initialization/jetty/ResponseIdentityHeaderHandler.java @@ -27,10 +27,14 @@ import org.eclipse.jetty.server.Handler; import org.eclipse.jetty.server.Request; import org.eclipse.jetty.util.Callback; +import javax.servlet.http.HttpServletRequest; import javax.servlet.http.HttpServletResponse; public class ResponseIdentityHeaderHandler extends Handler.Wrapper { + private static final String LOCAL_IDENTITY_ATTRIBUTE = + ResponseIdentityHeaderHandler.class.getName() + ".localIdentity"; + public static final String RESPONSE_SERVER_HEADER = "X-Druid-Server"; public static final String RESPONSE_SERVICE_HEADER = "X-Druid-Service"; public static final String RESPONSE_VERSION_HEADER = "X-Druid-Version"; @@ -98,6 +102,47 @@ public class ResponseIdentityHeaderHandler extends Handler.Wrapper headers.put(RESPONSE_VERSION_HEADER, responseVersion); } + public static void addIdentityHeaders( + final org.eclipse.jetty.server.Response response, + final DruidNode selfNode + ) + { + response.getHeaders().put(RESPONSE_SERVER_HEADER, selfNode.getHostAndPortToUse()); + response.getHeaders().put(RESPONSE_SERVICE_HEADER, selfNode.getServiceName()); + response.getHeaders().put(RESPONSE_VERSION_HEADER, selfNode.getVersion()); + } + + public static void rememberLocalIdentity( + final HttpServletRequest clientRequest, + final HttpServletResponse proxyResponse + ) + { + final String server = proxyResponse.getHeader(RESPONSE_SERVER_HEADER); + final String service = proxyResponse.getHeader(RESPONSE_SERVICE_HEADER); + final String version = proxyResponse.getHeader(RESPONSE_VERSION_HEADER); + if (server != null && service != null && version != null) { + clientRequest.setAttribute(LOCAL_IDENTITY_ATTRIBUTE, new ResponseIdentity(server, service, version)); + } + } + + public static void restoreLocalIdentity( + final HttpServletRequest clientRequest, + final HttpServletResponse proxyResponse + ) + { + final Object identity = clientRequest.getAttribute(LOCAL_IDENTITY_ATTRIBUTE); + if (identity instanceof ResponseIdentity localIdentity) { + clearRouterIdentity(proxyResponse); + proxyResponse.setHeader(RESPONSE_SERVER_HEADER, localIdentity.server()); + proxyResponse.setHeader(RESPONSE_SERVICE_HEADER, localIdentity.service()); + proxyResponse.setHeader(RESPONSE_VERSION_HEADER, localIdentity.version()); + } + } + + private record ResponseIdentity(String server, String service, String version) + { + } + public static void clearRouterIdentity(final HttpServletResponse proxyResponse) { // In EE8 compatible Jetty 12 using servlet API 4.x, setting a header to null is the accepted way to remove it. diff --git a/server/src/test/java/org/apache/druid/server/initialization/jetty/ResponseIdentityHeaderHandlerTest.java b/server/src/test/java/org/apache/druid/server/initialization/jetty/ResponseIdentityHeaderHandlerTest.java index 1040a3fc4db..d84cb40db50 100644 --- a/server/src/test/java/org/apache/druid/server/initialization/jetty/ResponseIdentityHeaderHandlerTest.java +++ b/server/src/test/java/org/apache/druid/server/initialization/jetty/ResponseIdentityHeaderHandlerTest.java @@ -29,7 +29,9 @@ import org.eclipse.jetty.server.Request; import org.eclipse.jetty.util.Callback; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import org.mockito.Mockito; +import javax.servlet.http.HttpServletRequest; import javax.servlet.http.HttpServletResponse; public class ResponseIdentityHeaderHandlerTest @@ -89,6 +91,31 @@ public class ResponseIdentityHeaderHandlerTest EasyMock.verify(proxyResponse); } + @Test + public void testRestoresRememberedLocalIdentity() + { + final HttpServletRequest clientRequest = Mockito.mock(HttpServletRequest.class); + final HttpServletResponse proxyResponse = Mockito.mock(HttpServletResponse.class); + final Object[] rememberedIdentity = new Object[1]; + Mockito.when(proxyResponse.getHeader(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER)) + .thenReturn("router:8888"); + Mockito.when(proxyResponse.getHeader(ResponseIdentityHeaderHandler.RESPONSE_SERVICE_HEADER)) + .thenReturn("druid/router"); + Mockito.when(proxyResponse.getHeader(ResponseIdentityHeaderHandler.RESPONSE_VERSION_HEADER)).thenReturn("39.0.0"); + Mockito.doAnswer(invocation -> { + rememberedIdentity[0] = invocation.getArgument(1); + return null; + }).when(clientRequest).setAttribute(Mockito.anyString(), Mockito.any()); + Mockito.when(clientRequest.getAttribute(Mockito.anyString())).thenAnswer(invocation -> rememberedIdentity[0]); + + ResponseIdentityHeaderHandler.rememberLocalIdentity(clientRequest, proxyResponse); + ResponseIdentityHeaderHandler.restoreLocalIdentity(clientRequest, proxyResponse); + + Mockito.verify(proxyResponse).setHeader(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER, "router:8888"); + Mockito.verify(proxyResponse).setHeader(ResponseIdentityHeaderHandler.RESPONSE_SERVICE_HEADER, "druid/router"); + Mockito.verify(proxyResponse).setHeader(ResponseIdentityHeaderHandler.RESPONSE_VERSION_HEADER, "39.0.0"); + } + @Test public void testShouldProxyIdentityHeaderWhenUpstreamReturnsAllHeaders() { diff --git a/services/src/main/java/org/apache/druid/server/AsyncQueryForwardingServlet.java b/services/src/main/java/org/apache/druid/server/AsyncQueryForwardingServlet.java index 0f54c3499b8..c26e8eb9af8 100644 --- a/services/src/main/java/org/apache/druid/server/AsyncQueryForwardingServlet.java +++ b/services/src/main/java/org/apache/druid/server/AsyncQueryForwardingServlet.java @@ -638,11 +638,25 @@ public class AsyncQueryForwardingServlet extends AsyncProxyServlet implements Qu // When response identity headers are enabled, the outer response handler initially adds the Router identity. // An upstream response must replace it with the upstream identity, or with no identity when the upstream does not // provide a complete header triple. + ResponseIdentityHeaderHandler.rememberLocalIdentity(clientRequest, proxyResponse); ResponseIdentityHeaderHandler.clearRouterIdentity(proxyResponse); StandardResponseHeaderFilterHolder.deduplicateHeadersInProxyServlet(proxyResponse, serverResponse); super.onServerResponseHeaders(clientRequest, proxyResponse, serverResponse); } + @Override + protected void onProxyResponseFailure( + final HttpServletRequest clientRequest, + final HttpServletResponse proxyResponse, + final Response serverResponse, + final Throwable failure + ) + { + // The error is generated by this Router, so replace any previously forwarded upstream identity with its own. + ResponseIdentityHeaderHandler.restoreLocalIdentity(clientRequest, proxyResponse); + super.onProxyResponseFailure(clientRequest, proxyResponse, serverResponse, failure); + } + @Override protected HttpField filterServerResponseHeader( HttpServletRequest clientRequest, --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
