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® 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¶m2=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]