This is an automated email from the ASF dual-hosted git repository.

FrankChen021 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git


The following commit(s) were added to refs/heads/master by this push:
     new 5e71056845b feat(server): add response identity headers (#20268)
5e71056845b is described below

commit 5e71056845b0a63fc01b09c4516f3d5fef1141a9
Author: Frank Chen <[email protected]>
AuthorDate: Thu Sep 10 14:06:49 2026 +0800

    feat(server): add response identity headers (#20268)
    
    * feat(server): add response identity headers
    
    * feat(server): include version in response identity
    
    * fix(server): preserve response identity across proxy paths
    
    * fix(server): identify locally generated proxy errors
    
    * Fix response identity test import order
---
 docs/api-reference/api-reference.md                |  16 ++
 docs/configuration/index.md                        |   8 +
 .../server/ResponseIdentityHeaderTest.java         | 259 +++++++++++++++++++++
 .../server/AsyncManagementForwardingServlet.java   |  41 +++-
 .../druid/server/http/OverlordProxyServlet.java    |  52 ++++-
 .../druid/server/initialization/ServerConfig.java  |  67 ++++++
 .../jetty/CliIndexerServerModule.java              |   1 +
 .../initialization/jetty/JettyServerModule.java    |  14 ++
 .../jetty/ResponseIdentityHeaderHandler.java       | 179 ++++++++++++++
 .../druid/initialization/ServerConfigTest.java     |   3 +
 .../AsyncManagementForwardingServletTest.java      |   3 +-
 .../server/http/OverlordProxyServletTest.java      |  92 +++++++-
 .../jetty/CliIndexerServerModuleTest.java          |  46 ++++
 .../jetty/ResponseIdentityHeaderHandlerTest.java   | 208 +++++++++++++++++
 .../druid/server/AsyncQueryForwardingServlet.java  |  27 +++
 15 files changed, 1012 insertions(+), 4 deletions(-)

diff --git a/docs/api-reference/api-reference.md 
b/docs/api-reference/api-reference.md
index d8fc86ee2a7..60ff4034660 100644
--- a/docs/api-reference/api-reference.md
+++ b/docs/api-reference/api-reference.md
@@ -26,6 +26,22 @@ sidebar_label: Overview
 
 This topic is an index to the Apache&circledR; Druid API documentation.
 
+Set `druid.server.http.enableResponseIdentityHeaders` to `true` to include 
headers on every HTTP response that identify
+the Druid service which generated the response:
+
+|Header|Description|
+|------|-----------|
+|`X-Druid-Server`|Advertised host and port of the Druid server that generated 
the response.|
+|`X-Druid-Service`|Service name of that Druid server, as configured by 
`druid.service`.|
+|`X-Druid-Version`|Version of that Druid server.|
+
+When a Router proxies a request, it passes through these headers from the 
upstream Druid server instead of returning
+the Router's identity. If the upstream server does not return all three 
headers, the Router does not return any of them.
+If the Router generates the response itself, the headers identify the Router.
+
+This feature is disabled by default because the headers may expose internal 
hostnames, IP addresses, ports, and the
+Druid version. Only enable it when clients are authorized to receive this 
information.
+
 ## HTTP APIs
 * [Druid SQL queries](./sql-api.md) to submit SQL queries using the Druid SQL 
API.
 * [SQL-based ingestion](./sql-ingestion-api.md) to submit SQL-based batch 
ingestion requests.
diff --git a/docs/configuration/index.md b/docs/configuration/index.md
index d07756b858e..14adf3150a3 100644
--- a/docs/configuration/index.md
+++ b/docs/configuration/index.md
@@ -556,6 +556,14 @@ Store task logs in HDFS. Note that the 
`druid-hdfs-storage` extension must be lo
 |`druid.indexer.logs.kill.initialDelay`| Optional. Number of milliseconds 
after Overlord start when first auto kill is run. |random value less than 
300000 (5 mins)|
 |`druid.indexer.logs.kill.delay`|Optional. Number of milliseconds of delay 
between successive executions of auto kill run. |21600000 (6 hours)|
 
+### Response identity headers
+
+This configuration applies to all Druid services.
+
+|Property|Description|Default|
+|--------|-----------|-------|
+|`druid.server.http.enableResponseIdentityHeaders`|If enabled, adds response 
headers containing the responding Druid service name, version, and advertised 
host and port. This can expose internal cluster topology and version 
information; only enable it when clients are authorized to receive this 
information.|`false`|
+
 ### API error response
 
 You can configure Druid API error responses to hide internal information like 
the Druid class name, stack trace, thread name, servlet name, code, line/column 
number, host, or IP address.
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
new file mode 100644
index 00000000000..333f8d32c30
--- /dev/null
+++ 
b/embedded-tests/src/test/java/org/apache/druid/testing/embedded/server/ResponseIdentityHeaderTest.java
@@ -0,0 +1,259 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.druid.testing.embedded.server;
+
+import 
org.apache.druid.server.initialization.jetty.ResponseIdentityHeaderHandler;
+import org.apache.druid.testing.embedded.EmbeddedBroker;
+import org.apache.druid.testing.embedded.EmbeddedCoordinator;
+import org.apache.druid.testing.embedded.EmbeddedDruidCluster;
+import org.apache.druid.testing.embedded.EmbeddedDruidServer;
+import org.apache.druid.testing.embedded.EmbeddedHistorical;
+import org.apache.druid.testing.embedded.EmbeddedOverlord;
+import org.apache.druid.testing.embedded.EmbeddedRouter;
+import org.apache.druid.testing.embedded.junit5.EmbeddedClusterTestBase;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.Timeout;
+
+import java.net.Socket;
+import java.net.URI;
+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;
+
+public class ResponseIdentityHeaderTest extends EmbeddedClusterTestBase
+{
+  private final EmbeddedCoordinator coordinator = new EmbeddedCoordinator();
+  private final EmbeddedOverlord overlord = new EmbeddedOverlord();
+  private final EmbeddedBroker broker = new EmbeddedBroker();
+  private final EmbeddedRouter router = new EmbeddedRouter();
+  private final EmbeddedHistorical historical = new EmbeddedHistorical();
+  private final HttpClient client = HttpClient.newHttpClient();
+
+  @Override
+  protected EmbeddedDruidCluster createCluster()
+  {
+    return EmbeddedDruidCluster.withEmbeddedDerbyAndZookeeper()
+                               
.addCommonProperty("druid.server.http.enableResponseIdentityHeaders", "true")
+                               .addServer(coordinator)
+                               .addServer(overlord)
+                               .addServer(broker)
+                               .addServer(historical)
+                               .addServer(router);
+  }
+
+  @AfterAll
+  public void closeClient()
+  {
+    client.close();
+  }
+
+  @Test
+  @Timeout(30)
+  public void testCoordinatorResponseIdentity_directAndViaRouter() throws 
Exception
+  {
+    assertResponseIdentity(
+        sendGet(getServerUrl(coordinator) + "/druid/coordinator/v1/isLeader"),
+        coordinator
+    );
+    assertResponseIdentity(
+        sendGet(getServerUrl(router) + "/druid/coordinator/v1/isLeader"),
+        coordinator
+    );
+  }
+
+  @Test
+  @Timeout(30)
+  public void testOverlordResponseIdentity_directViaCoordinatorAndViaRouter() 
throws Exception
+  {
+    assertResponseIdentity(
+        sendGet(getServerUrl(overlord) + "/druid/indexer/v1/isLeader"),
+        overlord
+    );
+    assertResponseIdentity(
+        sendGet(getServerUrl(coordinator) + "/druid/indexer/v1/isLeader"),
+        overlord
+    );
+    assertResponseIdentity(
+        sendGet(getServerUrl(router) + "/druid/indexer/v1/isLeader"),
+        overlord
+    );
+  }
+
+  @Test
+  @Timeout(30)
+  public void testBrokerNativeQueryResponseIdentity_directAndViaRouter() 
throws Exception
+  {
+    assertResponseIdentity(sendNativeQuery(broker), broker);
+    assertResponseIdentity(sendNativeQuery(router), broker);
+  }
+
+  @Test
+  @Timeout(30)
+  public void testBrokerSqlResponseIdentity_directAndViaRouter() throws 
Exception
+  {
+    assertResponseIdentity(sendSqlQuery(broker), broker);
+    assertResponseIdentity(sendSqlQuery(router), broker);
+  }
+
+  @Test
+  @Timeout(30)
+  public void testHistoricalNativeQueryResponseIdentity() throws Exception
+  {
+    assertResponseIdentity(sendNativeQuery(historical), historical);
+  }
+
+  @Test
+  @Timeout(30)
+  public void testRouterGeneratedAndEarlyErrorResponsesUseRouterIdentity() 
throws Exception
+  {
+    assertResponseIdentity(sendGet(getServerUrl(router) + "/status/health"), 
router);
+
+    final HttpRequest request = 
HttpRequest.newBuilder(URI.create(getServerUrl(router) + "/status/health"))
+                                           .timeout(Duration.ofSeconds(10))
+                                           .method("PATCH", 
HttpRequest.BodyPublishers.noBody())
+                                           .build();
+    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))
+                                           .timeout(Duration.ofSeconds(10))
+                                           .GET()
+                                           .build();
+    return client.send(request, HttpResponse.BodyHandlers.ofString());
+  }
+
+  private HttpResponse<String> sendNativeQuery(final EmbeddedDruidServer<?> 
server) throws Exception
+  {
+    final HttpRequest request = 
HttpRequest.newBuilder(URI.create(getServerUrl(server) + "/druid/v2"))
+                                           .header("Content-Type", 
"application/json")
+                                           .timeout(Duration.ofSeconds(10))
+                                           .POST(
+                                               
HttpRequest.BodyPublishers.ofString(
+                                                   """
+                                                   {
+                                                     "queryType": "timeseries",
+                                                     "dataSource": 
"missing_datasource",
+                                                     "granularity": "all",
+                                                     "intervals": 
["2000/3000"],
+                                                     "aggregations": [{"type": 
"count", "name": "rows"}]
+                                                   }
+                                                   """
+                                               )
+                                           )
+                                           .build();
+    return client.send(request, HttpResponse.BodyHandlers.ofString());
+  }
+
+  private HttpResponse<String> sendSqlQuery(final EmbeddedDruidServer<?> 
server) throws Exception
+  {
+    final HttpRequest request = 
HttpRequest.newBuilder(URI.create(getServerUrl(server) + "/druid/v2/sql"))
+                                           .header("Content-Type", 
"application/json")
+                                           .timeout(Duration.ofSeconds(10))
+                                           .POST(
+                                               
HttpRequest.BodyPublishers.ofString(
+                                                   """
+                                                   {
+                                                     "query": "SELECT 1"
+                                                   }
+                                                   """
+                                               )
+                                           )
+                                           .build();
+    return client.send(request, HttpResponse.BodyHandlers.ofString());
+  }
+
+  private static void assertResponseIdentity(
+      final HttpResponse<String> response,
+      final EmbeddedDruidServer<?> expectedServer
+  )
+  {
+    assertResponseIdentity(response, expectedServer, 200);
+  }
+
+  private static void assertResponseIdentity(
+      final HttpResponse<String> response,
+      final EmbeddedDruidServer<?> expectedServer,
+      final int expectedStatus
+  )
+  {
+    Assertions.assertEquals(expectedStatus, response.statusCode(), 
response.body());
+    Assertions.assertEquals(
+        List.of(expectedServer.bindings().selfNode().getHostAndPortToUse()),
+        
response.headers().allValues(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER)
+    );
+    Assertions.assertEquals(
+        List.of(expectedServer.bindings().selfNode().getServiceName()),
+        
response.headers().allValues(ResponseIdentityHeaderHandler.RESPONSE_SERVICE_HEADER)
+    );
+    Assertions.assertEquals(
+        List.of(expectedServer.bindings().selfNode().getVersion()),
+        
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 97b1a672118..042b0232fff 100644
--- 
a/server/src/main/java/org/apache/druid/server/AsyncManagementForwardingServlet.java
+++ 
b/server/src/main/java/org/apache/druid/server/AsyncManagementForwardingServlet.java
@@ -30,6 +30,8 @@ import org.apache.druid.guice.annotations.Global;
 import org.apache.druid.guice.annotations.Json;
 import org.apache.druid.guice.http.DruidHttpClientConfig;
 import org.apache.druid.java.util.common.StringUtils;
+import org.apache.druid.server.initialization.ServerConfig;
+import 
org.apache.druid.server.initialization.jetty.ResponseIdentityHeaderHandler;
 import 
org.apache.druid.server.initialization.jetty.StandardResponseHeaderFilterHolder;
 import org.apache.druid.server.security.AuthConfig;
 import org.apache.druid.server.security.AuthorizationUtils;
@@ -38,6 +40,7 @@ import org.eclipse.jetty.client.HttpClient;
 import org.eclipse.jetty.client.Request;
 import org.eclipse.jetty.client.Response;
 import org.eclipse.jetty.ee8.proxy.AsyncProxyServlet;
+import org.eclipse.jetty.http.HttpField;
 
 import javax.servlet.ServletException;
 import javax.servlet.http.HttpServletRequest;
@@ -75,6 +78,7 @@ public class AsyncManagementForwardingServlet extends 
AsyncProxyServlet
   private final DruidLeaderSelector coordLeaderSelector;
   private final DruidLeaderSelector overlordLeaderSelector;
   private final AuthorizerMapper authorizerMapper;
+  private final ServerConfig serverConfig;
 
   @Inject
   public AsyncManagementForwardingServlet(
@@ -83,7 +87,8 @@ public class AsyncManagementForwardingServlet extends 
AsyncProxyServlet
       @Global DruidHttpClientConfig httpClientConfig,
       @Coordinator DruidLeaderSelector coordLeaderSelector,
       @IndexingService DruidLeaderSelector overlordLeaderSelector,
-      AuthorizerMapper authorizerMapper
+      AuthorizerMapper authorizerMapper,
+      ServerConfig serverConfig
   )
   {
     this.jsonMapper = jsonMapper;
@@ -92,6 +97,7 @@ public class AsyncManagementForwardingServlet extends 
AsyncProxyServlet
     this.coordLeaderSelector = coordLeaderSelector;
     this.overlordLeaderSelector = overlordLeaderSelector;
     this.authorizerMapper = authorizerMapper;
+    this.serverConfig = serverConfig;
   }
 
   @Override
@@ -196,10 +202,43 @@ 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,
+      final Response serverResponse,
+      final HttpField field
+  )
+  {
+    // Identity headers are an all-or-nothing triple. Forward an identity 
field only when this Router has the feature
+    // enabled and the upstream response contains all three fields, so a 
client never observes a partial identity.
+    if (!ResponseIdentityHeaderHandler.shouldProxyIdentityHeader(
+        serverConfig.isEnableResponseIdentityHeaders(),
+        serverResponse,
+        field
+    )) {
+      return null;
+    }
+    return super.filterServerResponseHeader(clientRequest, serverResponse, 
field);
+  }
+
   /**
    * Authorizes router-internal requests that do not require any permissions. 
(But do require an authenticated user.)
    */
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 484d7858a0d..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
@@ -26,10 +26,15 @@ import org.apache.druid.guice.annotations.Global;
 import org.apache.druid.guice.http.DruidHttpClientConfig;
 import org.apache.druid.rpc.indexing.OverlordClient;
 import org.apache.druid.server.JettyUtils;
+import org.apache.druid.server.initialization.ServerConfig;
+import 
org.apache.druid.server.initialization.jetty.ResponseIdentityHeaderHandler;
+import 
org.apache.druid.server.initialization.jetty.StandardResponseHeaderFilterHolder;
 import org.apache.druid.server.security.AuthConfig;
 import org.eclipse.jetty.client.HttpClient;
 import org.eclipse.jetty.client.Request;
+import org.eclipse.jetty.client.Response;
 import org.eclipse.jetty.ee8.proxy.ProxyServlet;
+import org.eclipse.jetty.http.HttpField;
 
 import javax.servlet.ServletException;
 import javax.servlet.http.HttpServletRequest;
@@ -43,17 +48,20 @@ public class OverlordProxyServlet extends ProxyServlet
   private final OverlordClient overlordClient;
   private final Provider<HttpClient> httpClientProvider;
   private final DruidHttpClientConfig httpClientConfig;
+  private final ServerConfig serverConfig;
 
   @Inject
   OverlordProxyServlet(
       OverlordClient overlordClient,
       @Global Provider<HttpClient> httpClientProvider,
-      @Global DruidHttpClientConfig httpClientConfig
+      @Global DruidHttpClientConfig httpClientConfig,
+      ServerConfig serverConfig
   )
   {
     this.overlordClient = overlordClient;
     this.httpClientProvider = httpClientProvider;
     this.httpClientConfig = httpClientConfig;
+    this.serverConfig = serverConfig;
   }
 
   @Override
@@ -83,6 +91,48 @@ public class OverlordProxyServlet extends ProxyServlet
     return client;
   }
 
+  @Override
+  protected void onServerResponseHeaders(
+      final HttpServletRequest clientRequest,
+      final HttpServletResponse proxyResponse,
+      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,
+      final Response serverResponse,
+      final HttpField field
+  )
+  {
+    // The Coordinator-to-Overlord proxy follows the same all-or-nothing 
identity contract as Router proxies.
+    if (!ResponseIdentityHeaderHandler.shouldProxyIdentityHeader(
+        serverConfig.isEnableResponseIdentityHeaders(),
+        serverResponse,
+        field
+    )) {
+      return null;
+    }
+    return super.filterServerResponseHeader(clientRequest, serverResponse, 
field);
+  }
 
   @Override
   protected void sendProxyRequest(
diff --git 
a/server/src/main/java/org/apache/druid/server/initialization/ServerConfig.java 
b/server/src/main/java/org/apache/druid/server/initialization/ServerConfig.java
index 26082ab573a..093eae2d36f 100644
--- 
a/server/src/main/java/org/apache/druid/server/initialization/ServerConfig.java
+++ 
b/server/src/main/java/org/apache/druid/server/initialization/ServerConfig.java
@@ -88,6 +88,61 @@ public class ServerConfig
       @Nullable UriCompliance uriCompliance,
       boolean enforceStrictSNIHostChecking
   )
+  {
+    this(
+        numThreads,
+        queueSize,
+        enableRequestLimit,
+        maxIdleTime,
+        defaultQueryTimeout,
+        maxScatterGatherBytes,
+        maxSubqueryRows,
+        maxSubqueryBytes,
+        useNestedForUnknownTypeInSubquery,
+        maxQueryTimeout,
+        maxRequestHeaderSize,
+        gracefulShutdownTimeout,
+        unannouncePropagationDelay,
+        inflateBufferSize,
+        compressionLevel,
+        enableForwardedRequestCustomizer,
+        allowedHttpMethods,
+        showDetailedJettyErrors,
+        errorResponseTransformStrategy,
+        contentSecurityPolicy,
+        enableHSTS,
+        false,
+        uriCompliance,
+        enforceStrictSNIHostChecking
+    );
+  }
+
+  public ServerConfig(
+      int numThreads,
+      int queueSize,
+      boolean enableRequestLimit,
+      @NotNull Period maxIdleTime,
+      long defaultQueryTimeout,
+      long maxScatterGatherBytes,
+      int maxSubqueryRows,
+      String maxSubqueryBytes,
+      boolean useNestedForUnknownTypeInSubquery,
+      long maxQueryTimeout,
+      int maxRequestHeaderSize,
+      @NotNull Period gracefulShutdownTimeout,
+      @NotNull Period unannouncePropagationDelay,
+      int inflateBufferSize,
+      int compressionLevel,
+      boolean enableForwardedRequestCustomizer,
+      @NotNull List<String> allowedHttpMethods,
+      boolean showDetailedJettyErrors,
+      @NotNull ErrorResponseTransformStrategy errorResponseTransformStrategy,
+      @Nullable String contentSecurityPolicy,
+      boolean enableHSTS,
+      boolean enableResponseIdentityHeaders,
+      @Nullable UriCompliance uriCompliance,
+      boolean enforceStrictSNIHostChecking
+  )
   {
     this.numThreads = numThreads;
     this.queueSize = queueSize;
@@ -110,6 +165,7 @@ public class ServerConfig
     this.errorResponseTransformStrategy = errorResponseTransformStrategy;
     this.contentSecurityPolicy = contentSecurityPolicy;
     this.enableHSTS = enableHSTS;
+    this.enableResponseIdentityHeaders = enableResponseIdentityHeaders;
     this.uriCompliance = uriCompliance != null ? uriCompliance : 
UriCompliance.LEGACY;
     this.enforceStrictSNIHostChecking = enforceStrictSNIHostChecking;
   }
@@ -206,6 +262,9 @@ public class ServerConfig
   @JsonProperty
   private boolean enableHSTS = false;
 
+  @JsonProperty
+  private boolean enableResponseIdentityHeaders = false;
+
   @JsonProperty
   @JsonDeserialize(using = UriComplianceDeserializer.class)
   @JsonSerialize(using = UriComplianceSerializer.class)
@@ -330,6 +389,11 @@ public class ServerConfig
     return enableHSTS;
   }
 
+  public boolean isEnableResponseIdentityHeaders()
+  {
+    return enableResponseIdentityHeaders;
+  }
+
   public boolean isEnableQueryRequestsQueuing()
   {
     return enableQueryRequestsQueuing;
@@ -376,6 +440,7 @@ public class ServerConfig
            
errorResponseTransformStrategy.equals(that.errorResponseTransformStrategy) &&
            Objects.equals(contentSecurityPolicy, 
that.getContentSecurityPolicy()) &&
            enableHSTS == that.enableHSTS &&
+           enableResponseIdentityHeaders == that.enableResponseIdentityHeaders 
&&
            enableQueryRequestsQueuing == that.enableQueryRequestsQueuing &&
            Objects.equals(uriCompliance, that.uriCompliance) &&
            enforceStrictSNIHostChecking == that.enforceStrictSNIHostChecking;
@@ -406,6 +471,7 @@ public class ServerConfig
         showDetailedJettyErrors,
         contentSecurityPolicy,
         enableHSTS,
+        enableResponseIdentityHeaders,
         enableQueryRequestsQueuing,
         uriCompliance,
         enforceStrictSNIHostChecking
@@ -437,6 +503,7 @@ public class ServerConfig
            ", showDetailedJettyErrors=" + showDetailedJettyErrors +
            ", contentSecurityPolicy=" + contentSecurityPolicy +
            ", enableHSTS=" + enableHSTS +
+           ", enableResponseIdentityHeaders=" + enableResponseIdentityHeaders +
            ", enableQueryRequestsQueuing=" + enableQueryRequestsQueuing +
            ", uriCompliance=" + uriCompliance +
            ", enforceStrictSNIHostChecking=" + enforceStrictSNIHostChecking +
diff --git 
a/server/src/main/java/org/apache/druid/server/initialization/jetty/CliIndexerServerModule.java
 
b/server/src/main/java/org/apache/druid/server/initialization/jetty/CliIndexerServerModule.java
index f8768f3a7f0..9a7edd490ec 100644
--- 
a/server/src/main/java/org/apache/druid/server/initialization/jetty/CliIndexerServerModule.java
+++ 
b/server/src/main/java/org/apache/druid/server/initialization/jetty/CliIndexerServerModule.java
@@ -161,6 +161,7 @@ public class CliIndexerServerModule implements Module
         oldConfig.getErrorResponseTransformStrategy(),
         oldConfig.getContentSecurityPolicy(),
         oldConfig.isEnableHSTS(),
+        oldConfig.isEnableResponseIdentityHeaders(),
         oldConfig.getUriCompliance(),
         oldConfig.isEnforceStrictSNIHostChecking()
     );
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 9baa832ca4e..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
@@ -404,6 +404,9 @@ public class JettyServerModule extends JerseyServletModule
     JettyServerInitializer initializer = 
injector.getInstance(JettyServerInitializer.class);
     try {
       initializer.initialize(server, injector);
+      if (config.isEnableResponseIdentityHeaders()) {
+        server.setHandler(new ResponseIdentityHeaderHandler(node, 
server.getHandler()));
+      }
     }
     catch (Exception e) {
       throw new RE(e, "server initialization exception");
@@ -482,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
new file mode 100644
index 00000000000..a6e8a722014
--- /dev/null
+++ 
b/server/src/main/java/org/apache/druid/server/initialization/jetty/ResponseIdentityHeaderHandler.java
@@ -0,0 +1,179 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.druid.server.initialization.jetty;
+
+import org.apache.druid.server.DruidNode;
+import org.eclipse.jetty.client.Response;
+import org.eclipse.jetty.http.HttpField;
+import org.eclipse.jetty.http.HttpFields;
+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";
+
+  private final String responseServer;
+  private final String responseService;
+  private final String responseVersion;
+
+  public ResponseIdentityHeaderHandler(final DruidNode selfNode, final Handler 
handler)
+  {
+    super(handler);
+    responseServer = selfNode.getHostAndPortToUse();
+    responseService = selfNode.getServiceName();
+    responseVersion = selfNode.getVersion();
+  }
+
+  @Override
+  public boolean handle(
+      final Request request,
+      final org.eclipse.jetty.server.Response response,
+      final Callback callback
+  ) throws Exception
+  {
+    addIdentityHeaders(response);
+    final HttpFields.Mutable responseHeaders = new 
HttpFields.Mutable.Wrapper(response.getHeaders())
+    {
+      @Override
+      public HttpFields.Mutable clear()
+      {
+        super.clear();
+        addIdentityHeaders(this);
+        return this;
+      }
+    };
+    return super.handle(
+        request,
+        new org.eclipse.jetty.server.Response.Wrapper(request, response)
+        {
+          @Override
+          public HttpFields.Mutable getHeaders()
+          {
+            return responseHeaders;
+          }
+
+          @Override
+          public void reset()
+          {
+            super.reset();
+            addIdentityHeaders(this);
+          }
+        },
+        callback
+    );
+  }
+
+  private void addIdentityHeaders(final org.eclipse.jetty.server.Response 
response)
+  {
+    addIdentityHeaders(response.getHeaders());
+  }
+
+  private void addIdentityHeaders(final HttpFields.Mutable headers)
+  {
+    headers.put(RESPONSE_SERVER_HEADER, responseServer);
+    headers.put(RESPONSE_SERVICE_HEADER, responseService);
+    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.
+    proxyResponse.setHeader(RESPONSE_SERVER_HEADER, null);
+    proxyResponse.setHeader(RESPONSE_SERVICE_HEADER, null);
+    proxyResponse.setHeader(RESPONSE_VERSION_HEADER, null);
+  }
+
+  public static boolean shouldProxyIdentityHeader(
+      final boolean responseIdentityHeadersEnabled,
+      final Response serverResponse,
+      final HttpField field
+  )
+  {
+    if (!isIdentityHeader(field.getName())) {
+      return true;
+    }
+    if (!responseIdentityHeadersEnabled) {
+      return false;
+    }
+
+    final HttpFields headers = serverResponse.getHeaders();
+    return headers.contains(RESPONSE_SERVER_HEADER)
+           && headers.contains(RESPONSE_SERVICE_HEADER)
+           && headers.contains(RESPONSE_VERSION_HEADER);
+  }
+
+  private static boolean isIdentityHeader(final String headerName)
+  {
+    return RESPONSE_SERVER_HEADER.equalsIgnoreCase(headerName)
+           || RESPONSE_SERVICE_HEADER.equalsIgnoreCase(headerName)
+           || RESPONSE_VERSION_HEADER.equalsIgnoreCase(headerName);
+  }
+}
diff --git 
a/server/src/test/java/org/apache/druid/initialization/ServerConfigTest.java 
b/server/src/test/java/org/apache/druid/initialization/ServerConfigTest.java
index 44323cde3e6..463833b6f5e 100644
--- a/server/src/test/java/org/apache/druid/initialization/ServerConfigTest.java
+++ b/server/src/test/java/org/apache/druid/initialization/ServerConfigTest.java
@@ -44,6 +44,7 @@ public class ServerConfigTest
     Assertions.assertEquals(defaultConfig, defaultConfig2);
     
Assertions.assertFalse(defaultConfig2.isEnableForwardedRequestCustomizer());
     Assertions.assertFalse(defaultConfig2.isEnableHSTS());
+    Assertions.assertFalse(defaultConfig2.isEnableResponseIdentityHeaders());
     Assertions.assertEquals(UriCompliance.LEGACY, 
defaultConfig.getUriCompliance());
     Assertions.assertEquals(true, 
defaultConfig.isEnforceStrictSNIHostChecking());
 
@@ -69,6 +70,7 @@ public class ServerConfigTest
         new AllowedRegexErrorResponseTransformStrategy(ImmutableList.of(".*")),
         "my-cool-policy",
         true,
+        true,
         UriCompliance.RFC3986,
         false
     );
@@ -84,6 +86,7 @@ public class ServerConfigTest
     Assertions.assertEquals("my-cool-policy", 
modifiedConfig.getContentSecurityPolicy());
     Assertions.assertEquals("my-cool-policy", 
modifiedConfig2.getContentSecurityPolicy());
     Assertions.assertTrue(modifiedConfig2.isEnableHSTS());
+    Assertions.assertTrue(modifiedConfig2.isEnableResponseIdentityHeaders());
     Assertions.assertEquals(UriCompliance.RFC3986, 
modifiedConfig2.getUriCompliance());
     Assertions.assertFalse(modifiedConfig2.isEnforceStrictSNIHostChecking());
   }
diff --git 
a/server/src/test/java/org/apache/druid/server/AsyncManagementForwardingServletTest.java
 
b/server/src/test/java/org/apache/druid/server/AsyncManagementForwardingServletTest.java
index b16f68f1dc0..ef95647284e 100644
--- 
a/server/src/test/java/org/apache/druid/server/AsyncManagementForwardingServletTest.java
+++ 
b/server/src/test/java/org/apache/druid/server/AsyncManagementForwardingServletTest.java
@@ -523,7 +523,8 @@ public class AsyncManagementForwardingServletTest extends 
BaseJettyTest
               injector.getInstance(DruidHttpClientConfig.class),
               coordinatorLeaderSelector,
               overlordLeaderSelector,
-              new AuthorizerMapper(ImmutableMap.of("allowAll", new 
AllowAllAuthorizer(null)))
+              new AuthorizerMapper(ImmutableMap.of("allowAll", new 
AllowAllAuthorizer(null))),
+              injector.getInstance(ServerConfig.class)
           )
       );
 
diff --git 
a/server/src/test/java/org/apache/druid/server/http/OverlordProxyServletTest.java
 
b/server/src/test/java/org/apache/druid/server/http/OverlordProxyServletTest.java
index a0a78654ea2..f58e3afe134 100644
--- 
a/server/src/test/java/org/apache/druid/server/http/OverlordProxyServletTest.java
+++ 
b/server/src/test/java/org/apache/druid/server/http/OverlordProxyServletTest.java
@@ -21,7 +21,12 @@ package org.apache.druid.server.http;
 
 import com.google.common.util.concurrent.Futures;
 import org.apache.druid.rpc.indexing.OverlordClient;
+import org.apache.druid.server.initialization.ServerConfig;
+import 
org.apache.druid.server.initialization.jetty.ResponseIdentityHeaderHandler;
 import org.easymock.EasyMock;
+import org.eclipse.jetty.client.Response;
+import org.eclipse.jetty.http.HttpField;
+import org.eclipse.jetty.http.HttpFields;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
@@ -46,9 +51,94 @@ public class OverlordProxyServletTest
 
     EasyMock.replay(overlordClient, request);
 
-    URI uri = URI.create(new OverlordProxyServlet(overlordClient, null, 
null).rewriteTarget(request));
+    URI uri = URI.create(
+        new OverlordProxyServlet(overlordClient, null, null, new 
ServerConfig()).rewriteTarget(request)
+    );
     
Assertions.assertEquals("https://overlord:port/druid/over%3Alord/worker?param1=test&param2=test2";,
 uri.toString());
 
     EasyMock.verify(overlordClient, request);
   }
+
+  @Test
+  public void testProxyCompleteIdentityWhenEnabled()
+  {
+    final Response serverResponse = mockResponse(
+        HttpFields.from(
+            new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER, 
"overlord:8090"),
+            new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_SERVICE_HEADER, 
"druid/overlord"),
+            new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_VERSION_HEADER, "39.0.0")
+        ),
+        1
+    );
+    final HttpField serverField = new HttpField(
+        ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER,
+        "overlord:8090"
+    );
+
+    Assertions.assertEquals(
+        serverField,
+        new OverlordProxyServlet(null, null, null, enabledServerConfig())
+            .filterServerResponseHeader(null, serverResponse, serverField)
+    );
+    EasyMock.verify(serverResponse);
+  }
+
+  @Test
+  public void testSuppressPartialIdentity()
+  {
+    final Response serverResponse = mockResponse(
+        HttpFields.from(
+            new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER, 
"overlord:8090"),
+            new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_SERVICE_HEADER, 
"druid/overlord")
+        ),
+        1
+    );
+
+    Assertions.assertNull(
+        new OverlordProxyServlet(null, null, null, enabledServerConfig())
+            .filterServerResponseHeader(
+                null,
+                serverResponse,
+                new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER, "overlord:8090")
+            )
+    );
+    EasyMock.verify(serverResponse);
+  }
+
+  @Test
+  public void testSuppressIdentityWhenDisabled()
+  {
+    final Response serverResponse = EasyMock.strictMock(Response.class);
+    EasyMock.replay(serverResponse);
+
+    Assertions.assertNull(
+        new OverlordProxyServlet(null, null, null, new ServerConfig())
+            .filterServerResponseHeader(
+                null,
+                serverResponse,
+                new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER, "overlord:8090")
+            )
+    );
+    EasyMock.verify(serverResponse);
+  }
+
+  private static ServerConfig enabledServerConfig()
+  {
+    return new ServerConfig()
+    {
+      @Override
+      public boolean isEnableResponseIdentityHeaders()
+      {
+        return true;
+      }
+    };
+  }
+
+  private static Response mockResponse(final HttpFields headers, final int 
calls)
+  {
+    final Response serverResponse = EasyMock.strictMock(Response.class);
+    
EasyMock.expect(serverResponse.getHeaders()).andReturn(headers).times(calls);
+    EasyMock.replay(serverResponse);
+    return serverResponse;
+  }
 }
diff --git 
a/server/src/test/java/org/apache/druid/server/initialization/jetty/CliIndexerServerModuleTest.java
 
b/server/src/test/java/org/apache/druid/server/initialization/jetty/CliIndexerServerModuleTest.java
new file mode 100644
index 00000000000..0b2bcb23000
--- /dev/null
+++ 
b/server/src/test/java/org/apache/druid/server/initialization/jetty/CliIndexerServerModuleTest.java
@@ -0,0 +1,46 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.druid.server.initialization.jetty;
+
+import org.apache.druid.server.initialization.ServerConfig;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.util.Properties;
+
+public class CliIndexerServerModuleTest
+{
+  @Test
+  public void testAdjustedServerConfigPreservesResponseIdentityHeaders()
+  {
+    final ServerConfig oldConfig = new ServerConfig()
+    {
+      @Override
+      public boolean isEnableResponseIdentityHeaders()
+      {
+        return true;
+      }
+    };
+
+    final ServerConfig adjustedConfig = new CliIndexerServerModule(new 
Properties()).makeAdjustedServerConfig(oldConfig);
+
+    Assertions.assertTrue(adjustedConfig.isEnableResponseIdentityHeaders());
+  }
+}
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
new file mode 100644
index 00000000000..d84cb40db50
--- /dev/null
+++ 
b/server/src/test/java/org/apache/druid/server/initialization/jetty/ResponseIdentityHeaderHandlerTest.java
@@ -0,0 +1,208 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.druid.server.initialization.jetty;
+
+import org.apache.druid.server.DruidNode;
+import org.easymock.EasyMock;
+import org.eclipse.jetty.client.Response;
+import org.eclipse.jetty.http.HttpField;
+import org.eclipse.jetty.http.HttpFields;
+import org.eclipse.jetty.server.Handler;
+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
+{
+  @Test
+  public void testRestoresIdentityHeadersAfterResponseResetAndHeaderClear() 
throws Exception
+  {
+    final DruidNode node = new DruidNode("druid/test", "test-host", false, 
8080, null, true, false);
+    final Request request = EasyMock.strictMock(Request.class);
+    final org.eclipse.jetty.server.Response response = 
EasyMock.mock(org.eclipse.jetty.server.Response.class);
+    final HttpFields.Mutable headers = HttpFields.build();
+
+    EasyMock.expect(response.getHeaders()).andReturn(headers).times(2);
+    response.reset();
+    EasyMock.expectLastCall().andAnswer(
+        () -> {
+          headers.clear();
+          return null;
+        }
+    );
+    EasyMock.replay(request, response);
+
+    final Handler handler = new Handler.Abstract.NonBlocking()
+    {
+      @Override
+      public boolean handle(
+          final Request request,
+          final org.eclipse.jetty.server.Response response,
+          final Callback callback
+      )
+      {
+        response.getHeaders().clear();
+        assertIdentityHeaders(response.getHeaders(), node);
+        response.reset();
+        return true;
+      }
+    };
+
+    Assertions.assertTrue(new ResponseIdentityHeaderHandler(node, 
handler).handle(request, response, Callback.NOOP));
+    assertIdentityHeaders(headers, node);
+    EasyMock.verify(request, response);
+  }
+
+  @Test
+  public void testClearRouterIdentityRemovesAllHeaders()
+  {
+    final HttpServletResponse proxyResponse = 
EasyMock.strictMock(HttpServletResponse.class);
+    
proxyResponse.setHeader(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER, 
null);
+    EasyMock.expectLastCall().once();
+    
proxyResponse.setHeader(ResponseIdentityHeaderHandler.RESPONSE_SERVICE_HEADER, 
null);
+    EasyMock.expectLastCall().once();
+    
proxyResponse.setHeader(ResponseIdentityHeaderHandler.RESPONSE_VERSION_HEADER, 
null);
+    EasyMock.expectLastCall().once();
+
+    EasyMock.replay(proxyResponse);
+    ResponseIdentityHeaderHandler.clearRouterIdentity(proxyResponse);
+    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()
+  {
+    final Response serverResponse = mockResponse(
+        HttpFields.from(
+            new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER, 
"upstream:8082"),
+            new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_SERVICE_HEADER, 
"druid/broker"),
+            new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_VERSION_HEADER, "39.0.0")
+        ),
+        1
+    );
+
+    Assertions.assertTrue(
+        ResponseIdentityHeaderHandler.shouldProxyIdentityHeader(
+            true,
+            serverResponse,
+            new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER, "upstream:8082")
+        )
+    );
+    EasyMock.verify(serverResponse);
+  }
+
+  @Test
+  public void testShouldNotProxyPartialIdentity()
+  {
+    final Response serverResponse = mockResponse(
+        HttpFields.from(
+            new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER, 
"upstream:8082"),
+            new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_SERVICE_HEADER, "druid/broker")
+        ),
+        1
+    );
+
+    Assertions.assertFalse(
+        ResponseIdentityHeaderHandler.shouldProxyIdentityHeader(
+            true,
+            serverResponse,
+            new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER, "upstream:8082")
+        )
+    );
+    EasyMock.verify(serverResponse);
+  }
+
+  @Test
+  public void testShouldProxyUnrelatedHeader()
+  {
+    final Response serverResponse = EasyMock.strictMock(Response.class);
+    EasyMock.replay(serverResponse);
+
+    Assertions.assertTrue(
+        ResponseIdentityHeaderHandler.shouldProxyIdentityHeader(
+            false,
+            serverResponse,
+            new HttpField("Content-Type", "application/json")
+        )
+    );
+    EasyMock.verify(serverResponse);
+  }
+
+  @Test
+  public void testShouldNotProxyIdentityWhenRouterFeatureIsDisabled()
+  {
+    final Response serverResponse = EasyMock.strictMock(Response.class);
+    EasyMock.replay(serverResponse);
+
+    Assertions.assertFalse(
+        ResponseIdentityHeaderHandler.shouldProxyIdentityHeader(
+            false,
+            serverResponse,
+            new 
HttpField(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER, "upstream:8082")
+        )
+    );
+    EasyMock.verify(serverResponse);
+  }
+
+  private static Response mockResponse(final HttpFields headers, final int 
calls)
+  {
+    final Response serverResponse = EasyMock.strictMock(Response.class);
+    
EasyMock.expect(serverResponse.getHeaders()).andReturn(headers).times(calls);
+    EasyMock.replay(serverResponse);
+    return serverResponse;
+  }
+
+  private static void assertIdentityHeaders(final HttpFields headers, final 
DruidNode node)
+  {
+    Assertions.assertEquals("test-host:8080", 
headers.get(ResponseIdentityHeaderHandler.RESPONSE_SERVER_HEADER));
+    Assertions.assertEquals("druid/test", 
headers.get(ResponseIdentityHeaderHandler.RESPONSE_SERVICE_HEADER));
+    Assertions.assertEquals(node.getVersion(), 
headers.get(ResponseIdentityHeaderHandler.RESPONSE_VERSION_HEADER));
+  }
+}
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 7870e99ee60..c26e8eb9af8 100644
--- 
a/services/src/main/java/org/apache/druid/server/AsyncQueryForwardingServlet.java
+++ 
b/services/src/main/java/org/apache/druid/server/AsyncQueryForwardingServlet.java
@@ -48,6 +48,7 @@ import org.apache.druid.query.QueryMetrics;
 import org.apache.druid.query.QueryToolChestWarehouse;
 import org.apache.druid.server.initialization.ServerConfig;
 import org.apache.druid.server.initialization.jetty.HttpException;
+import 
org.apache.druid.server.initialization.jetty.ResponseIdentityHeaderHandler;
 import 
org.apache.druid.server.initialization.jetty.StandardResponseHeaderFilterHolder;
 import org.apache.druid.server.log.RequestLogger;
 import org.apache.druid.server.metrics.QueryCountStatsProvider;
@@ -634,10 +635,28 @@ public class AsyncQueryForwardingServlet extends 
AsyncProxyServlet implements Qu
     if (responseContext != null) {
       proxyResponse.setHeader(responseContext.getName(), 
responseContext.getValue());
     }
+    // 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,
@@ -648,6 +667,14 @@ public class AsyncQueryForwardingServlet extends 
AsyncProxyServlet implements Qu
     if 
(QueryResource.HEADER_RESPONSE_CONTEXT.equalsIgnoreCase(field.getName())) {
       return null;
     }
+    // Identity headers must be forwarded as a triple and only when this 
Router has the feature enabled.
+    if (!ResponseIdentityHeaderHandler.shouldProxyIdentityHeader(
+        serverConfig.isEnableResponseIdentityHeaders(),
+        serverResponse,
+        field
+    )) {
+      return null;
+    }
     return super.filterServerResponseHeader(clientRequest, serverResponse, 
field);
   }
 


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to