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]

Reply via email to