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 bd234ed1767 fix(druid): handle 429/503 HTML before JSON parse in 
DirectDruidClient (#20151)
bd234ed1767 is described below

commit bd234ed1767913b935d36acee6287538eab11763
Author: Jeremy Schoemaker <[email protected]>
AuthorDate: Sun Sep 20 21:01:21 2026 -0500

    fix(druid): handle 429/503 HTML before JSON parse in DirectDruidClient 
(#20151)
    
    * fix(druid): handle 429/503 HTML before JSON parse in DirectDruidClient
    
    Fix verified RED->GREEN. Broker masks 429/503 HTML as JsonParseException 
0x3c at DirectDruidClient.java:242
    
    * fix(druid): propagate handleResponse errors and detect HTML in chunked 
bodies
    
    NettyHttpClient discarded exceptions thrown from 
HttpResponseHandler#handleResponse
    by resolving the future to null before rethrowing, so callers never saw the 
failure.
    Now the future is failed with the original exception.
    
    DirectDruidClient's 429/503 and HTML detection only inspected the buffer 
attached to
    the initial HttpResponse, which is empty for chunked replies. handleChunk 
now retries
    the same first-non-whitespace-byte check against each chunk until the 
prefix resolves,
    so HTML delivered via chunked transfer is still caught before JSON parsing.
    
    Also fixes an ImportOrder violation from the prior commit.
    
    * Address review: only short-circuit non-JSON 429/503; keep handleChunk 
errors
    
    * fix(druid): address CodeQL findings in NettyHttpClientTest
    
    Replace the deprecated URL(String) constructor with URI.create(...).toURL()
    in the two new regression tests, and use String#isEmpty() instead of
    equals("") when skipping header lines in serveRawResponse.
    
    * fix(druid): honor Smile responses and preserve typed errors across chunk 
boundaries
    
    DirectDruidClient's non-JSON/HTML detection only ever accepted '{'/'[' as a
    structured-body marker, so a Smile-encoded 429/503 error body (what data
    servers actually return when queried in Smile, which is the production
    default) was misclassified as non-JSON and replaced with a synthesized
    QueryCapacityExceededException, discarding the server's real structured
    error. failIfNonJsonBody now branches on isSmile and accepts the Smile
    format header byte (0x3a) as well.
    
    Separately, when a later chunk's failIfNonJsonBody throws, handleResponse
    has already returned a finished response, so the transport's future is
    already resolved and the exception can only reach the caller through
    exceptionCaught. That path flattened the original typed QueryException
    (e.g. QueryCapacityExceededException) into a plain RE or, on the queued
    failure stream, an IOException wrapping it as a mere cause. The read path
    now rethrows a RuntimeException cause directly instead of re-wrapping it,
    and JsonParserIterator#convertException unwraps a QueryException nested
    one level as a cause so its concrete type and error code still survive.
    
    Also hardens bodyPreview: control characters (CR/LF included) are now
    stripped from the upstream body preview embedded in exception messages,
    closing off log/message line injection from an untrusted body regardless
    of how severe that channel actually is.
    
    * fix(druid): close failure-publication race, drop upstream body from error 
messages
    
    DirectDruidClient published a read failure as two separate atomics: fail
    (a message flag) written before failCause (its Throwable). A concurrent
    SequenceInputStream callback could observe fail set and enter
    failureException() before failCause's write became visible, falling back
    to a generic RE and losing a typed exception such as
    QueryCapacityExceededException thrown by failIfNonJsonBody from a later
    chunk. Both fields are now published as a single Failure object through
    one AtomicReference, so any reader that observes a non-null failure()
    always sees its cause too; there is no longer a window in which the two
    can be observed out of order. No other field pair in this class guards
    data behind a separately-published flag; fail/failCause was the only
    instance of that shape.
    
    Also drops the sanitized body preview from throwForNonJsonBody's
    exception messages entirely. The broker-to-data-server response is not a
    trusted boundary (a proxy in between can put diagnostics or reflected
    content in its error page), and the message is logged by
    JsonParserIterator and can reach the query error/trailer, so it should
    never carry upstream bytes regardless of how well they're sanitized. The
    message now reports bounded metadata instead: HTTP status,
    Content-Type when present, and body length.
    
    * fix(druid): normalize host for later-chunk QueryExceptions
    
    DirectDruidClient's rethrow of a later-chunk failure (e.g. a
    QueryCapacityExceededException from failIfNonJsonBody) preserves the
    concrete exception type but bypasses JsonParserIterator's host
    normalization: init() and next() only caught checked exceptions, so an
    unchecked QueryException thrown directly from the underlying
    InputStream's read() escaped both methods unconverted, carrying
    whatever host the data server gave it (QueryCapacityExceededException's
    withErrorMessageAndResolvedHost() uses the data server's own locally
    resolved hostname, or null) instead of DirectDruidClient's target host.
    
    Fixed by routing that case through the same convertException every
    other error path already uses. In next(), whose try block contains only
    real reads, a broad QueryException catch is enough. In init(), whose
    try block also throws QueryExceptions it has already converted itself
    (from timeoutQuery() and from two explicit convertException calls), a
    broad catch there would re-enter convertException on its own output,
    duplicating the warning log and rebuilding an equivalent object for no
    reason. So init() instead wraps only the two calls that can surface a
    raw, unconverted QueryException from the stream (parser construction,
    which can read ahead for encoding detection, and each nextToken() call)
    in small helper methods, leaving its own explicit converted throws
    outside that catch. One normalization rule, applied at every point a
    QueryException can reach this class, rather than two.
    
    Adds testLaterChunkQueryCapacityExceededReportsSameHostAsInitialResponse,
    asserting a later-chunk QueryCapacityExceededException and one arriving
    in a structured initial-response JSON body both report
    DirectDruidClient's configured host.
    
    * fix(druid): scope direct rethrow to QueryException, normalize error-body 
reads
    
    Two error-propagation gaps in the streamed response path.
    
    DirectDruidClient rethrew every RuntimeException recorded by
    setupResponseReadFailure as itself. That was meant to keep a later-chunk
    QueryCapacityExceededException from being flattened, but it also caught
    transport failures: a mid-stream disconnect arrives through exceptionCaught 
as
    Netty's ChannelException, which then reached the caller raw, without the
    "Query[id] url[url] failed with exception msg [...]" wrapper that is the 
only
    place the query id and target url appear. Both rethrow sites are now scoped 
to
    QueryException, so every other cause keeps its previous form.
    
    JsonParserIterator.init() called jp.getCodec().readValue(jp, 
QueryException.class)
    directly in its START_OBJECT branch. A structured error body can span 
chunks, so
    a later-chunk failure surfaces from inside that read rather than as the
    deserialized value, and escaped init() without convertException, reporting 
the
    data server's host instead of this iterator's. That read now goes through
    readStructuredError(), matching createParser() and readNextToken(); an 
exception
    off the stream is converted there and thrown, a successfully deserialized
    QueryException is returned for the caller to convert, so neither path 
converts
    twice.
    
    Regression tests for both, each verified to fail without its fix.
    
    * style(processing): put the netty imports where checkstyle wants them 🧹
    
    One import-order violation took down all eleven CI jobs on #20151. Every job
    builds druid-processing before it does anything else, so checkstyle failing 
in
    that module meant static-checks, strict-compilation, openrewrite, packaging,
    docker-tests and all seven unit-test shards died at the same line and never
    reached the server/ module where the actual change lives.
    
      NettyHttpClientTest.java:28:1: Wrong order for
      'io.netty.handler.codec.http.HttpContent' import. [ImportOrder]
    
    The three io.netty imports had landed below the org.apache.druid block. The
    ruleset is groups="*,javax,java" with ordered=true, so inside the first 
group
    io.netty sorts above org.apache. Moved them up, nothing else touched.
    
    The other four files this PR changes are clean. The javax -> java adjacency 
in
    DirectDruidClient and JsonParserIterator looks out of order but is 
suppressed on
    purpose by codestyle/checkstyle-suppressions.xml, and matches the rest of 
the
    tree.
    
    * fix(server): classify an empty 429/503 body as capacity exceeded at end 
of response 🕳️
    
    A proxy that answers 429/503 with no body, or only whitespace, never 
resolves
    the body prefix check: Netty skips handleChunk for an empty LastHttpContent 
and
    whitespace-only chunks return no prefix byte. done() then completed an empty
    stream that JsonParserIterator reported as a generic EOF instead of
    QueryCapacityExceededException.
    
    Finalize the unresolved prefix in done(): a 429/503 that never resolved 
throws
    the same typed exception the chunk path throws, before the stream completes,
    so it reaches the caller synchronously or through exceptionCaught exactly 
as a
    chunk-detected error page does. Successful empty responses are untouched.
    
    Tests: empty 503, whitespace-only 429, empty 503 routed through 
exceptionCaught
    after the initial response completed the future, and an empty 200 that still
    completes normally. The first two fail on 9947f370ba with "nothing was 
thrown".
    
    * fix(server): route a non-JSON 504 to QueryTimeoutException instead of the 
JSON parser 🕰️
    
    A gateway that times out is not a gateway that is out of capacity. Envoy's
    default answer for an upstream that is reachable but too slow is a 
plain-text
    504, and until now only 429/503 were classified, so that body went straight 
to
    the JSON parser and died as a JsonParseException on 'u'.
    
    504 now joins the proxy-error statuses but reports QueryTimeoutException 
rather
    than QueryCapacityExceededException — telling an operator to add data-server
    capacity for what is a latency problem sends them the wrong way. Both the
    body-prefix path and the empty-body end-of-response path handle it, and a 
504
    carrying a structured error body in the request's own format is still left 
to
    the normal JSON error path.
---
 .../java/util/http/client/NettyHttpClientTest.java | 250 +++++++++
 .../org/apache/druid/client/DirectDruidClient.java | 286 ++++++++++-
 .../apache/druid/client/JsonParserIterator.java    |  74 ++-
 .../apache/druid/client/DirectDruidClientTest.java | 565 ++++++++++++++++++++-
 .../druid/client/JsonParserIteratorTest.java       |  56 ++
 5 files changed, 1217 insertions(+), 14 deletions(-)

diff --git 
a/processing/src/test/java/org/apache/druid/java/util/http/client/NettyHttpClientTest.java
 
b/processing/src/test/java/org/apache/druid/java/util/http/client/NettyHttpClientTest.java
new file mode 100644
index 00000000000..745189356d4
--- /dev/null
+++ 
b/processing/src/test/java/org/apache/druid/java/util/http/client/NettyHttpClientTest.java
@@ -0,0 +1,250 @@
+/*
+ * 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.java.util.http.client;
+
+import com.google.common.util.concurrent.ListenableFuture;
+import com.google.common.util.concurrent.SettableFuture;
+import io.netty.handler.codec.http.HttpContent;
+import io.netty.handler.codec.http.HttpMethod;
+import io.netty.handler.codec.http.HttpResponse;
+import org.apache.druid.java.util.common.StringUtils;
+import org.apache.druid.java.util.common.lifecycle.Lifecycle;
+import org.apache.druid.java.util.http.client.response.ClientResponse;
+import org.apache.druid.java.util.http.client.response.HttpResponseHandler;
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.io.BufferedReader;
+import java.io.InputStreamReader;
+import java.io.OutputStream;
+import java.net.ServerSocket;
+import java.net.Socket;
+import java.net.URI;
+import java.nio.charset.StandardCharsets;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * Tests for {@link NettyHttpClient} exercising real socket I/O.
+ */
+public class NettyHttpClientTest
+{
+  /**
+   * A response handler whose {@link #handleResponse} throws, simulating a 
handler that rejects the response (for
+   * example because of an unexpected status code or content type) before any 
content has been processed.
+   */
+  private static class ThrowingResponseHandler implements 
HttpResponseHandler<Object, Object>
+  {
+    private final RuntimeException toThrow;
+
+    ThrowingResponseHandler(RuntimeException toThrow)
+    {
+      this.toThrow = toThrow;
+    }
+
+    @Override
+    public ClientResponse<Object> handleResponse(HttpResponse response, 
TrafficCop trafficCop)
+    {
+      throw toThrow;
+    }
+
+    @Override
+    public ClientResponse<Object> handleChunk(ClientResponse<Object> 
clientResponse, HttpContent chunk, long chunkNum)
+    {
+      return clientResponse;
+    }
+
+    @Override
+    public ClientResponse<Object> done(ClientResponse<Object> clientResponse)
+    {
+      return ClientResponse.finished(clientResponse.getObj());
+    }
+
+    @Override
+    public void exceptionCaught(ClientResponse<Object> clientResponse, 
Throwable e)
+    {
+      // Nothing to do.
+    }
+  }
+
+  /**
+   * Regression test: an exception thrown from {@link 
HttpResponseHandler#handleResponse} must fail the future
+   * returned by {@link HttpClient#go}, not resolve it to {@code null}. 
Previously, {@link NettyHttpClient} would
+   * call {@code retVal.set(null)} in its catch block before rethrowing, which 
meant the thrown exception was
+   * discarded and callers observed a successful null result instead of a 
failure.
+   */
+  @Test
+  public void testHandleResponseExceptionFailsFuture() throws Exception
+  {
+    final ExecutorService exec = Executors.newSingleThreadExecutor();
+    final ServerSocket serverSocket = new ServerSocket(0);
+    serveRawResponse(exec, serverSocket, "HTTP/1.1 200 OK\r\nContent-Length: 
2\r\n\r\n{}");
+
+    final Lifecycle lifecycle = new Lifecycle();
+    try {
+      final HttpClientConfig config = HttpClientConfig.builder().build();
+      final HttpClient client = HttpClientInit.createClient(config, lifecycle);
+
+      final RuntimeException expected = new RuntimeException("boom from 
handleResponse");
+      final ListenableFuture<Object> future = client.go(
+          new Request(
+              HttpMethod.GET,
+              URI.create(StringUtils.format("http://localhost:%d/";, 
serverSocket.getLocalPort())).toURL()
+          ),
+          new ThrowingResponseHandler(expected)
+      );
+
+      final ExecutionException e = 
Assertions.assertThrows(ExecutionException.class, future::get);
+      Assertions.assertSame(expected, e.getCause(), "the exception thrown by 
handleResponse must not be lost");
+    }
+    finally {
+      exec.shutdownNow();
+      serverSocket.close();
+      lifecycle.stop();
+    }
+  }
+
+  /**
+   * A response handler that (like DirectDruidClient) completes the future 
from {@link #handleResponse}, then
+   * throws from {@link #handleChunk} on a later chunk, and records whatever 
{@link #exceptionCaught} is eventually
+   * handed.
+   */
+  private static class ThrowingChunkHandler implements 
HttpResponseHandler<Object, Object>
+  {
+    private final RuntimeException toThrow;
+    private final SettableFuture<Throwable> caught = SettableFuture.create();
+
+    ThrowingChunkHandler(RuntimeException toThrow)
+    {
+      this.toThrow = toThrow;
+    }
+
+    @Override
+    public ClientResponse<Object> handleResponse(HttpResponse response, 
TrafficCop trafficCop)
+    {
+      return ClientResponse.finished("initial");
+    }
+
+    @Override
+    public ClientResponse<Object> handleChunk(ClientResponse<Object> 
clientResponse, HttpContent chunk, long chunkNum)
+    {
+      if (chunkNum >= 2) {
+        throw toThrow;
+      }
+      return clientResponse;
+    }
+
+    @Override
+    public ClientResponse<Object> done(ClientResponse<Object> clientResponse)
+    {
+      return ClientResponse.finished(clientResponse.getObj());
+    }
+
+    @Override
+    public void exceptionCaught(ClientResponse<Object> clientResponse, 
Throwable e)
+    {
+      caught.set(e);
+    }
+  }
+
+  /**
+   * Regression test: when the future has already been completed by {@link 
HttpResponseHandler#handleResponse} (as
+   * DirectDruidClient does for chunked responses), an exception thrown by 
{@link HttpResponseHandler#handleChunk}
+   * on a later chunk must be delivered to {@link 
HttpResponseHandler#exceptionCaught} as itself. Previously the
+   * catch block only closed the channel, so the handler instead saw the 
generic "Channel disconnected"
+   * {@link io.netty.channel.ChannelException} raised by the resulting 
disconnect, and the real cause was lost.
+   */
+  @Test
+  public void testHandleChunkExceptionReachesExceptionCaught() throws Exception
+  {
+    final ExecutorService exec = Executors.newSingleThreadExecutor();
+    final ServerSocket serverSocket = new ServerSocket(0);
+    serveRawResponse(
+        exec,
+        serverSocket,
+        "HTTP/1.1 200 OK\r\nTransfer-Encoding: 
chunked\r\n\r\n2\r\n{}\r\n6\r\n<html>\r\n0\r\n\r\n"
+    );
+
+    final Lifecycle lifecycle = new Lifecycle();
+    try {
+      final HttpClientConfig config = HttpClientConfig.builder().build();
+      final HttpClient client = HttpClientInit.createClient(config, lifecycle);
+
+      final RuntimeException expected = new RuntimeException("boom from 
handleChunk");
+      final ThrowingChunkHandler handler = new ThrowingChunkHandler(expected);
+      final ListenableFuture<Object> future = client.go(
+          new Request(
+              HttpMethod.GET,
+              URI.create(StringUtils.format("http://localhost:%d/";, 
serverSocket.getLocalPort())).toURL()
+          ),
+          handler
+      );
+
+      Assertions.assertEquals("initial", future.get(10, TimeUnit.SECONDS));
+      Assertions.assertSame(
+          expected,
+          handler.caught.get(10, TimeUnit.SECONDS),
+          "the exception thrown by handleChunk must reach exceptionCaught, not 
a generic channel-disconnected error"
+      );
+    }
+    finally {
+      exec.shutdownNow();
+      serverSocket.close();
+      lifecycle.stop();
+    }
+  }
+
+  /**
+   * Accepts connections on {@code serverSocket} until interrupted; for each, 
reads the request headers and writes
+   * {@code rawResponse} verbatim.
+   */
+  private static void serveRawResponse(ExecutorService exec, ServerSocket 
serverSocket, String rawResponse)
+  {
+    exec.submit(
+        new Runnable()
+        {
+          @Override
+          public void run()
+          {
+            while (!Thread.currentThread().isInterrupted()) {
+              try (
+                  Socket clientSocket = serverSocket.accept();
+                  BufferedReader in = new BufferedReader(
+                      new InputStreamReader(clientSocket.getInputStream(), 
StandardCharsets.UTF_8)
+                  );
+                  OutputStream out = clientSocket.getOutputStream()
+              ) {
+                while (!in.readLine().isEmpty()) {
+                  // skip lines
+                }
+                out.write(rawResponse.getBytes(StandardCharsets.UTF_8));
+                out.flush();
+              }
+              catch (Exception e) {
+                // Suppress
+              }
+            }
+          }
+        }
+    );
+  }
+}
diff --git 
a/server/src/main/java/org/apache/druid/client/DirectDruidClient.java 
b/server/src/main/java/org/apache/druid/client/DirectDruidClient.java
index 5b4e921e505..b222c9bcc9c 100644
--- a/server/src/main/java/org/apache/druid/client/DirectDruidClient.java
+++ b/server/src/main/java/org/apache/druid/client/DirectDruidClient.java
@@ -21,6 +21,7 @@ package org.apache.druid.client;
 
 import com.fasterxml.jackson.databind.JavaType;
 import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.dataformat.smile.SmileConstants;
 import com.fasterxml.jackson.dataformat.smile.SmileFactory;
 import com.fasterxml.jackson.jaxrs.smile.SmileMediaTypes;
 import com.google.common.base.Preconditions;
@@ -49,7 +50,10 @@ import 
org.apache.druid.java.util.http.client.response.StatusResponseHandler;
 import org.apache.druid.java.util.http.client.response.StatusResponseHolder;
 import org.apache.druid.query.Queries;
 import org.apache.druid.query.Query;
+import org.apache.druid.query.QueryCapacityExceededException;
 import org.apache.druid.query.QueryContext;
+import org.apache.druid.query.QueryException;
+import org.apache.druid.query.QueryInterruptedException;
 import org.apache.druid.query.QueryMetrics;
 import org.apache.druid.query.QueryPlus;
 import org.apache.druid.query.QueryRunner;
@@ -65,6 +69,7 @@ import org.apache.druid.server.QueryResource;
 import org.apache.druid.utils.CloseableUtils;
 import org.joda.time.Duration;
 
+import javax.annotation.Nullable;
 import javax.ws.rs.core.MediaType;
 import java.io.IOException;
 import java.io.InputStream;
@@ -184,7 +189,39 @@ public class DirectDruidClient<T> implements QueryRunner<T>
         // handleResponse). Once set, incoming chunks are dropped rather than 
buffered.
         private final AtomicBoolean discard = new AtomicBoolean(false);
         private final AtomicBoolean nodeMetricsEmitted = new 
AtomicBoolean(false);
-        private final AtomicReference<String> fail = new AtomicReference<>();
+        // Tracks whether the response body's leading (non-whitespace) byte 
has been classified as JSON or
+        // HTML. For chunked responses the initial HttpResponse can arrive 
with an empty or all-whitespace body,
+        // in which case the check is retried against each subsequent 
HttpChunk until it resolves.
+        private final AtomicBoolean bodyPrefixResolved = new 
AtomicBoolean(false);
+        // HTTP status of the initial response, remembered so that an HTML 
body first seen in a later chunk of a
+        // 429/503 reply is still reported as capacity-exceeded rather than as 
a generic HTML-instead-of-JSON error.
+        private volatile int responseStatusCode = -1;
+        // Netty 4 delivers the body after the headers, so the Content-Type 
recorded in handleResponse has
+        // to survive until the first chunk arrives to be usable in the 
non-JSON check.
+        private volatile String responseContentType = null;
+        // Message and (when available) original exception for a read failure, 
published together through a single
+        // reference so a reader that observes a non-null Failure also sees 
its cause. Publishing them as two
+        // separate AtomicReferences (a "fail" flag written before a 
"failCause" detail) would let a concurrent
+        // SequenceInputStream callback observe the flag before the cause and 
fall back to a generic RE, losing a
+        // typed exception like QueryCapacityExceededException thrown by 
failIfNonJsonBody from a later chunk.
+        private final AtomicReference<Failure> failure = new 
AtomicReference<>();
+
+        /**
+         * Immutable pairing of a failure message with its originating 
exception, when one exists. Published as a
+         * single unit through {@link #failure} so the message and cause are 
always visible together.
+         */
+        private static final class Failure
+        {
+          private final String message;
+          private final Throwable cause;
+
+          private Failure(String message, Throwable cause)
+          {
+            this.message = message;
+            this.cause = cause;
+          }
+        }
+
         private final AtomicReference<TrafficCop> trafficCopRef = new 
AtomicReference<>();
 
         private QueryMetrics<? super Query<T>> queryMetrics;
@@ -237,12 +274,167 @@ public class DirectDruidClient<T> implements 
QueryRunner<T>
           return holder.getStream();
         }
 
+        /**
+         * Scans past leading whitespace in {@code buffer} looking for the 
first content byte, without consuming
+         * (advancing the reader index of) the buffer, and returns it. Once a 
non-whitespace byte is found, the prefix
+         * is considered resolved (see {@link #bodyPrefixResolved}) and later 
calls return null. If {@code buffer} is
+         * empty or entirely whitespace, the prefix remains unresolved (and 
null is returned) so a later call, from a
+         * subsequent chunk, can retry the check; this matters for chunked 
responses, where the initial
+         * {@link HttpResponse} can carry an empty body and the real content, 
HTML or otherwise, only arrives via
+         * {@link #handleChunk}.
+         */
+        @Nullable
+        private Byte bodyPrefixByte(ByteBuf buffer)
+        {
+          if (bodyPrefixResolved.get()) {
+            return null;
+          }
+          final int readerIndex = buffer.readerIndex();
+          final int readable = buffer.readableBytes();
+          for (int i = 0; i < readable; i++) {
+            byte b = buffer.getByte(readerIndex + i);
+            if (b == ' ' || b == '\n' || b == '\r' || b == '\t') {
+              continue;
+            }
+            bodyPrefixResolved.set(true);
+            return b;
+          }
+          return null;
+        }
+
+        /**
+         * Classifies the body prefix in {@code buffer} (see {@link 
#bodyPrefixByte}) and fails the query if it is not
+         * JSON (or Smile, when the request was sent as Smile per {@link 
#isSmile}). HTML always fails; any other
+         * non-JSON/non-Smile body fails only when the status is one of the 
proxy error statuses in
+         * {@link #isProxyErrorStatus}, since Druid itself never sends such a 
body with those statuses but proxies
+         * routinely do (an HTML error page from nginx or a load balancer, a 
plain-text "upstream connect error" or
+         * "upstream request timeout" from Envoy). A structured body in the 
request's own format, whatever the
+         * status, is left alone so that the normal parse path can surface the 
server's own structured error.
+         *
+         * @param contentType Content-Type header of the initial response, 
possibly null; a text/html value fails the
+         *                    query regardless of the body prefix
+         * @param chunkNum    0 for the initial response body, else the chunk 
number
+         */
+        private void failIfNonJsonBody(String contentType, ByteBuf buffer, 
long chunkNum)
+        {
+          final boolean isHtmlContentType =
+              contentType != null && 
StringUtils.toLowerCase(contentType).contains("text/html");
+          final Byte prefix = bodyPrefixByte(buffer);
+          final boolean isHtml = isHtmlContentType || (prefix != null && 
prefix == '<');
+          // A data server negotiates its response format from the request 
(ResourceIOReaderWriterFactory#factorize),
+          // so a Smile request gets a Smile response, error bodies included; 
those begin with the Smile format
+          // header's 0x3a byte rather than JSON's '{'/'['. Checking only 
'{'/'[' here would misclassify every
+          // structured Smile 429/503 body as non-JSON and discard the 
server's real error.
+          final boolean isNonJson = isSmile
+                                     ? prefix != null && prefix != 
SmileConstants.HEADER_BYTE_1
+                                     : prefix != null && prefix != '{' && 
prefix != '[';
+          final int statusCode = responseStatusCode;
+          if (isHtml || (isNonJson && isProxyErrorStatus(statusCode))) {
+            throwForNonJsonBody(statusCode, contentType, buffer, chunkNum, 
isHtml);
+          }
+        }
+
+        /**
+         * Whether {@code statusCode} is one an intermediary produces on the 
broker-to-data-server hop with a body of
+         * its own rather than Druid's: 429 and 503 when the upstream is 
rejecting or unreachable, 504 when it is
+         * reachable but did not answer in time. A 504 is not a capacity 
signal, so {@link #throwForNonJsonBody}
+         * reports it as a {@link QueryTimeoutException} while 429/503 stay 
{@link QueryCapacityExceededException}.
+         */
+        private boolean isProxyErrorStatus(int statusCode)
+        {
+          return statusCode == 429 || statusCode == 503 || statusCode == 504;
+        }
+
+        /**
+         * Fails the query because the response body is not JSON, typically an 
error page produced by a load balancer
+         * or reverse proxy sitting in front of the data server. A 429/503 
status is reported as
+         * {@link QueryCapacityExceededException} since that is what such 
intermediaries return when the server is
+         * over capacity; a 504 is reported as a {@link 
QueryTimeoutException}, since a gateway timeout says the
+         * upstream was reachable but did not answer in time, which is the 
same condition the client-side
+         * {@link #checkQueryTimeout} reports and not a capacity rejection; 
any other status is reported as a
+         * {@link QueryInterruptedException}. Either way, the caller gets a 
message that says what actually came back
+         * instead of a {@code JsonParseException} on {@code '<'}.
+         * <p>
+         * The exception message never includes the body itself. The 
broker-to-data-server response is not a trusted
+         * boundary (it can be an error page from any proxy sitting in 
between), and this message is logged by
+         * {@link JsonParserIterator} and can reach the query error/trailer, 
so it is limited to bounded, sanitized
+         * metadata: HTTP status, Content-Type when present, and the body 
length.
+         *
+         * @param statusCode  HTTP status of the initial response
+         * @param contentType Content-Type header of the initial response, 
possibly null
+         * @param buffer      the buffer in which the non-JSON body was 
detected (the initial response body or a chunk)
+         * @param chunkNum    0 if detected in the initial response body, else 
the chunk number
+         * @param isHtml      whether the body was identified as HTML 
specifically (vs. some other non-JSON content)
+         */
+        private void throwForNonJsonBody(
+            int statusCode,
+            String contentType,
+            ByteBuf buffer,
+            long chunkNum,
+            boolean isHtml
+        )
+        {
+          final String where = chunkNum > 0 ? StringUtils.format(" (detected 
in chunk[%d])", chunkNum) : "";
+          final String bodyInfo = StringUtils.format(
+              "contentType[%s] bodyLength[%d]",
+              contentType,
+              buffer.readableBytes()
+          );
+          if (statusCode == 504) {
+            throw new QueryTimeoutException(
+                StringUtils.format(
+                    "Query[%s] url[%s] timed out upstream with status[%s]%s 
%s",
+                    query.getId(),
+                    url,
+                    statusCode,
+                    where,
+                    bodyInfo
+                ),
+                host
+            );
+          }
+          if (statusCode == 429 || statusCode == 503) {
+            throw 
QueryCapacityExceededException.withErrorMessageAndResolvedHost(
+                StringUtils.format(
+                    "Query[%s] url[%s] failed with status[%s]%s %s",
+                    query.getId(),
+                    url,
+                    statusCode,
+                    where,
+                    bodyInfo
+                )
+            );
+          }
+          throw new QueryInterruptedException(
+              QueryException.UNKNOWN_EXCEPTION_ERROR_CODE,
+              StringUtils.format(
+                  "Query[%s] url[%s] returned %s response instead of JSON with 
status[%s]%s %s",
+                  query.getId(),
+                  url,
+                  isHtml ? "HTML" : "non-JSON",
+                  statusCode,
+                  where,
+                  bodyInfo
+              ),
+              QueryInterruptedException.class.getName(),
+              host
+          );
+        }
+
         @Override
         public ClientResponse<InputStream> handleResponse(HttpResponse 
response, TrafficCop trafficCop)
         {
           trafficCopRef.set(trafficCop);
           checkQueryTimeout();
-          // Netty 4: initial HttpResponse has no body content; body arrives 
as HttpContent chunks.
+          // Netty 4: the initial HttpResponse carries no body, so the status 
and Content-Type are recorded
+          // here and the body itself is inspected on the first HttpContent 
chunk. The goal is to detect a
+          // non-JSON body (a 429/503/504 error page from a proxy, say) before 
it reaches the JSON parser, where it
+          // would surface only as a JsonParseException on 0x3c ('<'). The 
shortcut is taken only when the body
+          // is confirmed non-JSON: a 429/503 carrying Druid's own JSON error 
body (a genuine
+          // QueryCapacityExceededException or SERVICE_UNAVAILABLE from a data 
server) falls through to the
+          // normal JSON error path below, which preserves the server's 
structured error details.
+          responseStatusCode = response.status().code();
+          responseContentType = 
response.headers().get(HttpHeaders.Names.CONTENT_TYPE);
 
           log.debug("Initial response from url[%s] for queryId[%s]", url, 
query.getId());
           responseStartTimeNs = System.nanoTime();
@@ -299,8 +491,8 @@ public class DirectDruidClient<T> implements QueryRunner<T>
                       if (discard.get()) {
                         return false;
                       }
-                      if (fail.get() != null) {
-                        throw new RE(fail.get());
+                      if (failure.get() != null) {
+                        throw failureException();
                       }
                       checkQueryTimeout();
 
@@ -314,8 +506,8 @@ public class DirectDruidClient<T> implements QueryRunner<T>
                     @Override
                     public InputStream nextElement()
                     {
-                      if (fail.get() != null) {
-                        throw new RE(fail.get());
+                      if (failure.get() != null) {
+                        throw failureException();
                       }
 
                       try {
@@ -372,6 +564,13 @@ public class DirectDruidClient<T> implements QueryRunner<T>
 
           checkTotalBytesLimit(bytes);
 
+          // Under Netty 4 the body never appears on the initial HttpResponse, 
so this is where the JSON-vs-not
+          // prefix check happens, against the first chunk that carries any 
non-whitespace byte, and before those
+          // bytes are enqueued for JSON parsing. Otherwise an HTML error page 
(from a load balancer or reverse
+          // proxy, say) is enqueued blind and surfaces later as a confusing 
JsonParseException. The Content-Type
+          // comes from the headers recorded in handleResponse. This is a 
no-op once the prefix has resolved.
+          failIfNonJsonBody(responseContentType, channelBuffer, chunkNum);
+
           boolean continueReading = true;
           if (bytes > 0) {
             try {
@@ -388,9 +587,51 @@ public class DirectDruidClient<T> implements QueryRunner<T>
           return ClientResponse.finished(clientResponse.getObj(), 
continueReading);
         }
 
+        /**
+         * Classifies a proxy error status (see {@link #isProxyErrorStatus}) 
whose body never resolved the prefix
+         * check, that is, a body that was empty or all whitespace. Netty 
skips {@link #handleChunk} for an empty
+         * {@code LastHttpContent} and whitespace-only chunks leave {@link 
#bodyPrefixResolved} unset, so end of
+         * response is the last chance to report such a reply as {@link 
QueryCapacityExceededException} (429/503) or
+         * {@link QueryTimeoutException} (504); otherwise the empty stream 
completes normally and
+         * {@link JsonParserIterator} reports it as a generic EOF. A 
successful empty response is left alone.
+         */
+        private void failIfUnresolvedProxyErrorBody()
+        {
+          final int statusCode = responseStatusCode;
+          if (bodyPrefixResolved.get() || !isProxyErrorStatus(statusCode)) {
+            return;
+          }
+          if (statusCode == 504) {
+            throw new QueryTimeoutException(
+                StringUtils.format(
+                    "Query[%s] url[%s] timed out upstream with status[%s] and 
no JSON body contentType[%s] bodyLength[%d]",
+                    query.getId(),
+                    url,
+                    statusCode,
+                    responseContentType,
+                    totalByteCount.get()
+                ),
+                host
+            );
+          }
+          throw QueryCapacityExceededException.withErrorMessageAndResolvedHost(
+              StringUtils.format(
+                  "Query[%s] url[%s] failed with status[%s] and no JSON body 
contentType[%s] bodyLength[%d]",
+                  query.getId(),
+                  url,
+                  statusCode,
+                  responseContentType,
+                  totalByteCount.get()
+              )
+          );
+        }
+
         @Override
         public ClientResponse<InputStream> done(ClientResponse<InputStream> 
clientResponse)
         {
+          // Runs before the stream is completed so the failure reaches the 
caller the same way a chunk-detected
+          // non-JSON body does: synchronously here, or via exceptionCaught in 
NettyHttpClient.
+          failIfUnresolvedProxyErrorBody();
           long stopTimeNs = System.nanoTime();
           long nodeTimeNs = stopTimeNs - requestStartTimeNs;
           final long nodeTimeMs = TimeUnit.NANOSECONDS.toMillis(nodeTimeNs);
@@ -442,7 +683,9 @@ public class DirectDruidClient<T> implements QueryRunner<T>
         private void setupResponseReadFailure(String msg, Throwable th)
         {
           emitNodeMetrics(System.nanoTime() - requestStartTimeNs);
-          fail.set(msg);
+          // Publish message and cause together as one Failure so a reader 
that observes a non-null failure()
+          // always sees the cause that goes with it; see the field comment on 
failure.
+          failure.set(new Failure(msg, th));
           queue.clear();
           queue.offer(
               InputStreamHolder.fromStream(
@@ -451,7 +694,15 @@ public class DirectDruidClient<T> implements QueryRunner<T>
                     @Override
                     public int read() throws IOException
                     {
-                      if (th != null) {
+                      if (th instanceof QueryException) {
+                        // Rethrow a typed query failure (e.g. 
QueryCapacityExceededException) as itself rather than
+                        // burying it as the cause of a generic IOException, 
where it would otherwise only be
+                        // recoverable by callers that specifically unwrap 
getCause(). Deliberately limited to
+                        // QueryException: a transport-level RuntimeException 
such as Netty's ChannelException from a
+                        // mid-stream disconnect must keep going out as an 
IOException, because that is the form
+                        // JsonParserIterator normalizes into a 
QueryInterruptedException carrying this client's host.
+                        throw (QueryException) th;
+                      } else if (th != null) {
                         throw new IOException(msg, th);
                       } else {
                         throw new IOException(msg);
@@ -464,6 +715,25 @@ public class DirectDruidClient<T> implements QueryRunner<T>
           );
         }
 
+        /**
+         * Returns the exception to surface for a failure recorded by {@link 
#setupResponseReadFailure}. Rethrows the
+         * original cause directly when it is already a {@link QueryException} 
(e.g. the
+         * {@link QueryCapacityExceededException} thrown from {@link 
#handleChunk} on a later chunk) so its concrete
+         * type survives to the caller instead of being flattened into a plain 
{@link RE}. Every other cause keeps
+         * the pre-existing {@link RE} form, so a transport failure such as a 
mid-stream disconnect still reaches
+         * the caller with this method's message rather than as a raw Netty 
exception. Only called after
+         * confirming {@link #failure} is non-null, so the message and cause 
it reads are always the ones from the
+         * same {@link Failure} publication.
+         */
+        private RuntimeException failureException()
+        {
+          final Failure f = failure.get();
+          if (f.cause instanceof QueryException) {
+            return (QueryException) f.cause;
+          }
+          return new RE(f.message);
+        }
+
         // Emit exactly once, regardless of whether we reach this via done() 
or setupResponseReadFailure().
         private void emitNodeMetrics(long nodeTimeNs)
         {
diff --git 
a/server/src/main/java/org/apache/druid/client/JsonParserIterator.java 
b/server/src/main/java/org/apache/druid/client/JsonParserIterator.java
index 7aa88774397..3c86c0d6313 100644
--- a/server/src/main/java/org/apache/druid/client/JsonParserIterator.java
+++ b/server/src/main/java/org/apache/druid/client/JsonParserIterator.java
@@ -132,6 +132,14 @@ public class JsonParserIterator<T> implements 
CloseableIterator<T>
         throw convertException(e);
       }
     }
+    catch (QueryException e) {
+      // A QueryException can reach here unwrapped, straight from a 
InputStream.read() call inside jp.nextToken()/
+      // readValue() above, when DirectDruidClient rethrows a later-chunk 
failure (e.g. QueryCapacityExceededException
+      // from failIfNonJsonBody) as itself rather than as the cause of an 
IOException. Route it through the same
+      // convertException as every other error path so it gets this iterator's 
host instead of whatever host it
+      // happened to carry when the data server (rather than the broker) 
constructed it.
+      throw convertException(e);
+    }
   }
 
   @Override
@@ -173,7 +181,7 @@ public class JsonParserIterator<T> implements 
CloseableIterator<T>
         InputStream is = hasTimeout ? future.get(timeLeftMillis, 
TimeUnit.MILLISECONDS) : future.get();
 
         if (is != null) {
-          jp = objectMapper.getFactory().createParser(is);
+          jp = createParser(is);
         } else if (checkTimeout()) {
           throw timeoutQuery();
         } else {
@@ -188,12 +196,12 @@ public class JsonParserIterator<T> implements 
CloseableIterator<T>
           );
         }
 
-        final JsonToken nextToken = jp.nextToken();
+        final JsonToken nextToken = readNextToken();
         if (nextToken == JsonToken.START_ARRAY) {
-          jp.nextToken();
+          readNextToken();
           objectCodec = jp.getCodec();
         } else if (nextToken == JsonToken.START_OBJECT) {
-          throw convertException(jp.getCodec().readValue(jp, 
QueryException.class));
+          throw convertException(readStructuredError());
         } else {
           String errMsg = jp.getValueAsString();
           if (errMsg != null) {
@@ -221,6 +229,56 @@ public class JsonParserIterator<T> implements 
CloseableIterator<T>
     }
   }
 
+  /**
+   * Creates the parser for {@code is}, converting a {@link QueryException} 
that reaches here unwrapped (e.g. a
+   * later-chunk {@link QueryCapacityExceededException} DirectDruidClient 
rethrows as itself rather than as the
+   * cause of an {@link IOException}, since Jackson's stream bootstrapping can 
read ahead for encoding detection
+   * before any token is parsed) through {@link #convertException}, same as 
every other error path. Scoped to just
+   * this call, rather than caught around the whole {@link #init()} try block, 
so it never re-catches the
+   * QueryExceptions {@link #init()} throws explicitly after already calling 
{@link #convertException}.
+   */
+  private JsonParser createParser(InputStream is) throws IOException
+  {
+    try {
+      return objectMapper.getFactory().createParser(is);
+    }
+    catch (QueryException e) {
+      throw convertException(e);
+    }
+  }
+
+  /**
+   * Deserializes the structured error body {@link #init()} found in place of 
a result array, with the same
+   * unwrapped-{@link QueryException} handling as {@link #createParser}: the 
body can span chunks, so a later-chunk
+   * failure can surface from inside this read rather than from the 
deserialized value. An exception that comes off
+   * the stream is converted here and thrown; a QueryException that was 
successfully deserialized is returned for the
+   * caller to convert, so neither path is converted twice.
+   */
+  private QueryException readStructuredError() throws IOException
+  {
+    try {
+      return jp.getCodec().readValue(jp, QueryException.class);
+    }
+    catch (QueryException e) {
+      throw convertException(e);
+    }
+  }
+
+  /**
+   * Reads the next {@link JsonToken} from {@link #jp}, with the same 
unwrapped-{@link QueryException} handling as
+   * {@link #createParser}, for the same reason: a later-chunk failure can 
surface on any read from the underlying
+   * stream, not only the first.
+   */
+  private JsonToken readNextToken() throws IOException
+  {
+    try {
+      return jp.nextToken();
+    }
+    catch (QueryException e) {
+      throw convertException(e);
+    }
+  }
+
   private QueryTimeoutException timeoutQuery()
   {
     return new QueryTimeoutException(StringUtils.nonStrictFormat("url[%s] 
timed out", url), host);
@@ -238,6 +296,14 @@ public class JsonParserIterator<T> implements 
CloseableIterator<T>
   private QueryException convertException(Throwable cause)
   {
     LOG.warn(cause, "Query [%s] to host [%s] interrupted", queryId, host);
+    // A QueryException thrown on the transport thread (e.g. 
QueryCapacityExceededException from a later chunk in
+    // DirectDruidClient) can reach here re-wrapped as the cause of a plain 
RE/IOException rather than as itself, if
+    // it had to cross the boundary between the Netty callback that detected 
it and the thread reading the response.
+    // Unwrap that one level so its concrete type and error code still make it 
out instead of falling through to a
+    // generic QueryInterruptedException below.
+    if (!(cause instanceof QueryException) && cause.getCause() instanceof 
QueryException) {
+      cause = cause.getCause();
+    }
     if (cause instanceof QueryException) {
       final QueryException queryException = (QueryException) cause;
       if (queryException.getErrorCode() == null) {
diff --git 
a/server/src/test/java/org/apache/druid/client/DirectDruidClientTest.java 
b/server/src/test/java/org/apache/druid/client/DirectDruidClientTest.java
index c0f39e3b5a1..ca50b052ffd 100644
--- a/server/src/test/java/org/apache/druid/client/DirectDruidClientTest.java
+++ b/server/src/test/java/org/apache/druid/client/DirectDruidClientTest.java
@@ -20,24 +20,40 @@
 package org.apache.druid.client;
 
 import com.fasterxml.jackson.databind.ObjectMapper;
+import com.fasterxml.jackson.dataformat.smile.SmileFactory;
+import com.fasterxml.jackson.jaxrs.smile.SmileMediaTypes;
 import com.google.common.util.concurrent.Futures;
 import com.google.common.util.concurrent.ListenableFuture;
 import com.google.common.util.concurrent.SettableFuture;
+import io.netty.buffer.Unpooled;
+import io.netty.channel.ChannelException;
+import io.netty.handler.codec.http.DefaultHttpContent;
+import io.netty.handler.codec.http.DefaultHttpResponse;
+import io.netty.handler.codec.http.HttpContent;
+import io.netty.handler.codec.http.HttpHeaders;
 import io.netty.handler.codec.http.HttpMethod;
+import io.netty.handler.codec.http.HttpResponse;
+import io.netty.handler.codec.http.HttpResponseStatus;
+import io.netty.handler.codec.http.HttpVersion;
 import io.netty.handler.timeout.ReadTimeoutException;
 import org.apache.druid.data.input.ResourceInputSource;
 import org.apache.druid.jackson.DefaultObjectMapper;
 import org.apache.druid.java.util.common.DateTimes;
 import org.apache.druid.java.util.common.ISE;
+import org.apache.druid.java.util.common.RE;
 import org.apache.druid.java.util.common.StringUtils;
 import org.apache.druid.java.util.common.guava.Sequence;
 import org.apache.druid.java.util.emitter.service.ServiceEmitter;
 import org.apache.druid.java.util.http.client.HttpClient;
 import org.apache.druid.java.util.http.client.Request;
+import org.apache.druid.java.util.http.client.response.ClientResponse;
+import org.apache.druid.java.util.http.client.response.HttpResponseHandler;
 import org.apache.druid.java.util.metrics.StubServiceEmitter;
 import org.apache.druid.query.Druids;
 import org.apache.druid.query.NestedDataTestUtils;
+import org.apache.druid.query.QueryCapacityExceededException;
 import org.apache.druid.query.QueryContexts;
+import org.apache.druid.query.QueryException;
 import org.apache.druid.query.QueryInterruptedException;
 import org.apache.druid.query.QueryPlus;
 import org.apache.druid.query.QueryRunnerTestHelper;
@@ -58,6 +74,7 @@ import org.apache.druid.server.metrics.NoopServiceEmitter;
 import org.apache.druid.testing.TemporaryFolderExtension;
 import org.apache.druid.timeline.DataSegment;
 import org.apache.druid.timeline.SegmentId;
+import org.joda.time.Duration;
 import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.BeforeEach;
@@ -71,6 +88,7 @@ import java.io.PipedInputStream;
 import java.io.PipedOutputStream;
 import java.net.MalformedURLException;
 import java.net.URL;
+import java.util.Arrays;
 import java.util.List;
 import java.util.Map;
 import java.util.concurrent.CancellationException;
@@ -400,17 +418,560 @@ public class DirectDruidClientTest
         actualException.getMessage());
   }
 
+  @Test
+  public void testHtml503InInitialResponseIsCapacityExceeded()
+  {
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(HttpResponseStatus.SERVICE_UNAVAILABLE, 
"text/html", "<html><body>503 Service Unavailable</body></html>")
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryCapacityExceededException e = Assertions.assertThrows(
+        QueryCapacityExceededException.class,
+        () -> client.run(queryPlus, responseContext)
+    );
+    Assertions.assertTrue(e.getMessage().contains("status[503]"), 
e.getMessage());
+    Assertions.assertTrue(e.getMessage().contains("contentType[text/html]"), 
e.getMessage());
+    Assertions.assertFalse(e.getMessage().contains("<html>"), e.getMessage());
+  }
+
+  @Test
+  public void testHtml503InLaterChunkAfterEmptyInitialBodyIsCapacityExceeded()
+  {
+    // A chunked 503 whose initial HttpResponse carries an empty body: the 
HTML only shows up in a later chunk, and
+    // must still be classified as capacity-exceeded rather than reaching the 
JSON parser or being reported as a
+    // generic HTML-instead-of-JSON error.
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(HttpResponseStatus.SERVICE_UNAVAILABLE, null, 
"", "  \n", "<html><body>503</body></html>")
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryCapacityExceededException e = Assertions.assertThrows(
+        QueryCapacityExceededException.class,
+        () -> client.run(queryPlus, responseContext)
+    );
+    Assertions.assertTrue(e.getMessage().contains("status[503]"), 
e.getMessage());
+    Assertions.assertTrue(e.getMessage().contains("detected in chunk[2]"), 
e.getMessage());
+  }
+
+  @Test
+  public void testEmptyBody503IsCapacityExceeded()
+  {
+    // A proxy answering 503 with no body at all. Netty skips handleChunk for 
an empty LastHttpContent, so the body
+    // prefix never resolves and done() is the only place left to classify the 
response; without that it completes
+    // an empty stream that JsonParserIterator later reports as a generic EOF 
instead of capacity-exceeded.
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(HttpResponseStatus.SERVICE_UNAVAILABLE, null)
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryCapacityExceededException e = Assertions.assertThrows(
+        QueryCapacityExceededException.class,
+        () -> client.run(queryPlus, responseContext)
+    );
+    Assertions.assertTrue(e.getMessage().contains("status[503]"), 
e.getMessage());
+    Assertions.assertTrue(e.getMessage().contains("bodyLength[0]"), 
e.getMessage());
+  }
+
+  @Test
+  public void testWhitespaceOnlyBody429IsCapacityExceeded()
+  {
+    // Whitespace-only chunks never resolve the prefix either, so the same 
end-of-response classification applies.
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(HttpResponseStatus.TOO_MANY_REQUESTS, 
"text/plain", "  \n", "\t")
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryCapacityExceededException e = Assertions.assertThrows(
+        QueryCapacityExceededException.class,
+        () -> client.run(queryPlus, responseContext)
+    );
+    Assertions.assertTrue(e.getMessage().contains("status[429]"), 
e.getMessage());
+    Assertions.assertTrue(e.getMessage().contains("contentType[text/plain]"), 
e.getMessage());
+    Assertions.assertTrue(e.getMessage().contains("bodyLength[4]"), 
e.getMessage());
+  }
+
+  @Test
+  public void testEmptyBody503RoutedThroughExceptionCaughtIsCapacityExceeded()
+  {
+    // Same lifecycle as 
testHtml503InLaterChunkAfterFinishedInitialResponseIsCapacityExceeded: the 
future is already
+    // completed by handleResponse, so the end-of-response classification can 
only reach the caller through
+    // exceptionCaught and the vended stream, and must keep its type on the 
way.
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(HttpResponseStatus.SERVICE_UNAVAILABLE, null, 
true)
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryCapacityExceededException e = Assertions.assertThrows(
+        QueryCapacityExceededException.class,
+        () -> client.run(queryPlus, responseContext).toList()
+    );
+    Assertions.assertTrue(e.getMessage().contains("status[503]"), 
e.getMessage());
+    Assertions.assertTrue(e.getMessage().contains("bodyLength[0]"), 
e.getMessage());
+  }
+
+  @Test
+  public void testEmptyBody200IsNotShortCircuited()
+  {
+    // A successful response with no body keeps completing normally; only 
429/503 are classified at end of response.
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(HttpResponseStatus.OK, null)
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    Assertions.assertDoesNotThrow(() -> client.run(queryPlus, 
responseContext));
+  }
+
+  @Test
+  public void testHtmlInLaterChunkOf200ResponseIsQueryInterrupted()
+  {
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(HttpResponseStatus.OK, null, "", 
"<html><body>oops</body></html>")
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryInterruptedException e = Assertions.assertThrows(
+        QueryInterruptedException.class,
+        () -> client.run(queryPlus, responseContext)
+    );
+    Assertions.assertTrue(e.getMessage().contains("returned HTML response 
instead of JSON"), e.getMessage());
+    Assertions.assertTrue(e.getMessage().contains("detected in chunk[1]"), 
e.getMessage());
+  }
+
+  @Test
+  public void testPlainText503IsCapacityExceeded()
+  {
+    // Not every proxy error page is HTML: Envoy, for one, returns a 
plain-text body with a 503. The body itself is
+    // never echoed into the message (see 
testNonJsonBodyMessageContainsNoRawBodyBytes); this checks the status and
+    // Content-Type metadata that stands in for it.
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(HttpResponseStatus.SERVICE_UNAVAILABLE, 
"text/plain", "upstream connect error or disconnect/reset before headers")
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryCapacityExceededException e = Assertions.assertThrows(
+        QueryCapacityExceededException.class,
+        () -> client.run(queryPlus, responseContext)
+    );
+    Assertions.assertTrue(e.getMessage().contains("status[503]"), 
e.getMessage());
+    Assertions.assertTrue(e.getMessage().contains("contentType[text/plain]"), 
e.getMessage());
+    Assertions.assertFalse(e.getMessage().contains("upstream connect error"), 
e.getMessage());
+  }
+
+  @Test
+  public void testPlainText504IsQueryTimeout()
+  {
+    // Envoy's default answer for an upstream that is reachable but too slow 
is a plain-text 504, not a 503. That is
+    // a timeout, not a capacity rejection, so it must surface as 
QueryTimeoutException: reporting it as
+    // capacity-exceeded would tell an operator to add data-server capacity 
for what is a latency problem, and
+    // leaving it unclassified sends the plain-text body to the JSON parser to 
die as a JsonParseException on 'u'.
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(HttpResponseStatus.GATEWAY_TIMEOUT, 
"text/plain", "upstream request timeout")
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryTimeoutException e = Assertions.assertThrows(
+        QueryTimeoutException.class,
+        () -> client.run(queryPlus, responseContext)
+    );
+    Assertions.assertTrue(e.getMessage().contains("status[504]"), 
e.getMessage());
+    Assertions.assertTrue(e.getMessage().contains("contentType[text/plain]"), 
e.getMessage());
+    Assertions.assertFalse(e.getMessage().contains("upstream request 
timeout"), e.getMessage());
+  }
+
+  @Test
+  public void testHtml504IsQueryTimeout()
+  {
+    // An HTML 504 from nginx takes the same route. HTML fails the body check 
on any status, so this pins which
+    // exception the HTML path produces for 504 rather than whether it fails 
at all.
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(HttpResponseStatus.GATEWAY_TIMEOUT, 
"text/html", "<html><body>504 Gateway Time-out</body></html>")
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryTimeoutException e = Assertions.assertThrows(
+        QueryTimeoutException.class,
+        () -> client.run(queryPlus, responseContext)
+    );
+    Assertions.assertTrue(e.getMessage().contains("status[504]"), 
e.getMessage());
+  }
+
+  @Test
+  public void testEmptyBody504IsQueryTimeout()
+  {
+    // The end-of-response path classifies 504 too: a gateway that times out 
with no body at all would otherwise
+    // complete an empty stream and reach the caller as a generic EOF.
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(HttpResponseStatus.GATEWAY_TIMEOUT, null)
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryTimeoutException e = Assertions.assertThrows(
+        QueryTimeoutException.class,
+        () -> client.run(queryPlus, responseContext)
+    );
+    Assertions.assertTrue(e.getMessage().contains("status[504]"), 
e.getMessage());
+    Assertions.assertTrue(e.getMessage().contains("bodyLength[0]"), 
e.getMessage());
+  }
+
+  @Test
+  public void testJson504IsNotShortCircuited()
+  {
+    // Same rule as 429/503: a 504 carrying a structured error body in the 
request's own format is left to the
+    // normal JSON error path so the server's own message survives.
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(
+            HttpResponseStatus.GATEWAY_TIMEOUT,
+            "application/json",
+            "{\"error\":\"Unknown exception\",\"errorMessage\":\"backend says 
no\",\"errorClass\":\"x\",\"host\":\"h\"}"
+        )
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryException e = Assertions.assertThrows(
+        QueryException.class,
+        () -> client.run(queryPlus, responseContext).toList()
+    );
+    Assertions.assertEquals("backend says no", e.getMessage());
+  }
+
+  @Test
+  public void testNonJsonBodyMessageContainsNoRawBodyBytes()
+  {
+    // The exception message must never echo the upstream body itself: the 
broker-to-data-server response is not a
+    // trusted boundary, and this message is logged by JsonParserIterator and 
can reach the query error/trailer. It
+    // is limited to bounded sanitized metadata instead: HTTP status, 
Content-Type, and body length.
+    final String body = "upstream down\r\nX-Injected: evil\nsecond line";
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(
+            HttpResponseStatus.SERVICE_UNAVAILABLE,
+            "text/plain",
+            body
+        )
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryCapacityExceededException e = Assertions.assertThrows(
+        QueryCapacityExceededException.class,
+        () -> client.run(queryPlus, responseContext)
+    );
+    Assertions.assertFalse(e.getMessage().contains("upstream down"), 
e.getMessage());
+    Assertions.assertFalse(e.getMessage().contains("X-Injected"), 
e.getMessage());
+    Assertions.assertFalse(e.getMessage().contains("second line"), 
e.getMessage());
+    Assertions.assertFalse(e.getMessage().contains("\r"), e.getMessage());
+    Assertions.assertFalse(e.getMessage().contains("\n"), e.getMessage());
+    Assertions.assertTrue(e.getMessage().contains("status[503]"), 
e.getMessage());
+    Assertions.assertTrue(e.getMessage().contains("contentType[text/plain]"), 
e.getMessage());
+    Assertions.assertTrue(
+        e.getMessage().contains("bodyLength[" + 
StringUtils.toUtf8(body).length + "]"),
+        e.getMessage()
+    );
+  }
+
+  @Test
+  public void testJson503IsNotShortCircuited()
+  {
+    // A 503 carrying Druid's own JSON error body must take the normal JSON 
error path so the server's message
+    // survives, instead of being replaced by a synthesized capacity-exceeded 
error.
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(
+            HttpResponseStatus.SERVICE_UNAVAILABLE,
+            "application/json",
+            "{\"error\":\"Unknown exception\",\"errorMessage\":\"backend says 
no\",\"errorClass\":\"x\",\"host\":\"h\"}"
+        )
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryException e = Assertions.assertThrows(
+        QueryException.class,
+        () -> client.run(queryPlus, responseContext).toList()
+    );
+    Assertions.assertFalse(e instanceof QueryCapacityExceededException, 
e.getClass().getName());
+    Assertions.assertEquals("backend says no", e.getMessage());
+  }
+
+  @Test
+  public void testSmile503IsNotShortCircuited() throws IOException
+  {
+    // Same as testJson503IsNotShortCircuited, but for the Smile ObjectMapper 
DirectDruidClientFactory actually
+    // injects in production. A Smile-encoded structured error body starts 
with the Smile format header byte (0x3a),
+    // not '{'/'[', and must not be misclassified as a non-JSON proxy error 
page.
+    final ObjectMapper smileObjectMapper = new DefaultObjectMapper(new 
SmileFactory(), null);
+    final byte[] smileBody = smileObjectMapper.writeValueAsBytes(
+        new QueryException("Unknown exception", "backend says no", "x", "h")
+    );
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(
+            HttpResponseStatus.SERVICE_UNAVAILABLE,
+            SmileMediaTypes.APPLICATION_JACKSON_SMILE,
+            false,
+            smileBody
+        ),
+        smileObjectMapper
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryException e = Assertions.assertThrows(
+        QueryException.class,
+        () -> client.run(queryPlus, responseContext).toList()
+    );
+    Assertions.assertFalse(e instanceof QueryCapacityExceededException, 
e.getClass().getName());
+    Assertions.assertEquals("backend says no", e.getMessage());
+  }
+
+  @Test
+  public void 
testHtml503InLaterChunkAfterFinishedInitialResponseIsCapacityExceeded()
+  {
+    // Reproduces the real NettyHttpClient lifecycle instead of the simplified 
one the other later-chunk tests use:
+    // handleResponse always returns an already-finished ClientResponse, so 
the transport's future is completed
+    // before any chunk is seen, and a later chunk's exception can only reach 
the caller through exceptionCaught.
+    // Before the fix, that path re-wrapped the QueryCapacityExceededException 
thrown by handleChunk into a plain
+    // RE (or, via the queued failure InputStream, an IOException), losing the 
original type entirely.
+    final DirectDruidClient client = makeDirectDruidClient(
+        new ScriptedHttpClient(
+            HttpResponseStatus.SERVICE_UNAVAILABLE,
+            null,
+            true,
+            StringUtils.toUtf8(""),
+            StringUtils.toUtf8("<html><body>503</body></html>")
+        )
+    );
+
+    final QueryPlus queryPlus = getQueryPlus();
+    final QueryCapacityExceededException e = Assertions.assertThrows(
+        QueryCapacityExceededException.class,
+        () -> client.run(queryPlus, responseContext).toList()
+    );
+    Assertions.assertTrue(e.getMessage().contains("status[503]"), 
e.getMessage());
+    Assertions.assertTrue(e.getMessage().contains("detected in chunk[1]"), 
e.getMessage());
+  }
+
+  @Test
+  public void 
testLaterChunkQueryCapacityExceededReportsSameHostAsInitialResponse()
+  {
+    // DirectDruidClient rethrows a later-chunk QueryCapacityExceededException 
as itself (see the previous test), but
+    // it constructs that exception with 
QueryCapacityExceededException.withErrorMessageAndResolvedHost(), whose host
+    // is the data server's own locally resolved hostname (or null). An 
initial-response QueryCapacityExceededException
+    // arriving as Druid's own structured JSON error body (the case 
testJson503IsNotShortCircuited also exercises, as
+    // opposed to the failIfNonJsonBody HTML/non-JSON shortcut, which throws 
synchronously out of handleResponse and
+    // never reaches JsonParserIterator at all) is deserialized by 
JsonParserIterator.init()'s START_OBJECT branch and
+    // normalized by convertException to DirectDruidClient.host. Unless the 
later-chunk exception is also routed
+    // through convertException, the two report different hosts for the same 
logical failure. This asserts they
+    // match, both landing on DirectDruidClient's own configured host 
(hostName), not the data server's.
+    final DirectDruidClient initialResponseClient = makeDirectDruidClient(
+        new ScriptedHttpClient(
+            HttpResponseStatus.SERVICE_UNAVAILABLE,
+            "application/json",
+            "{\"error\":\"Query capacity exceeded\",\"errorMessage\":\"too 
many queries\","
+            + 
"\"errorClass\":\"org.apache.druid.query.QueryCapacityExceededException\",\"host\":\"data-server-01:8100\"}"
+        )
+    );
+    final QueryCapacityExceededException initialResponseException = 
Assertions.assertThrows(
+        QueryCapacityExceededException.class,
+        () -> initialResponseClient.run(getQueryPlus(), 
responseContext).toList()
+    );
+
+    final DirectDruidClient laterChunkClient = makeDirectDruidClient(
+        new ScriptedHttpClient(
+            HttpResponseStatus.SERVICE_UNAVAILABLE,
+            null,
+            true,
+            StringUtils.toUtf8(""),
+            StringUtils.toUtf8("<html><body>503</body></html>")
+        )
+    );
+    final QueryCapacityExceededException laterChunkException = 
Assertions.assertThrows(
+        QueryCapacityExceededException.class,
+        () -> laterChunkClient.run(getQueryPlus(), responseContext).toList()
+    );
+
+    Assertions.assertEquals(hostName, initialResponseException.getHost());
+    Assertions.assertEquals(hostName, laterChunkException.getHost());
+    Assertions.assertEquals(initialResponseException.getHost(), 
laterChunkException.getHost());
+  }
+
+  @Test
+  public void testMidStreamTransportFailureIsNotRethrownRaw()
+  {
+    // A mid-stream disconnect is not a query error: it reaches 
DirectDruidClient through exceptionCaught as a plain
+    // transport RuntimeException (Netty's ChannelException here), after 
handleResponse has already vended the stream.
+    // Rethrowing the cause is scoped to QueryException precisely so this case 
keeps the pre-existing behaviour of
+    // surfacing the "Query[id] url[url] failed with exception msg [...]" 
wrapper, which is the only place the query
+    // id and the target url appear at all, instead of a bare ChannelException 
that carries neither.
+    final DirectDruidClient client = makeDirectDruidClient(
+        new TransportFailureHttpClient(new ChannelException("connection reset 
by peer"))
+    );
+
+    final RE e = Assertions.assertThrows(
+        RE.class,
+        () -> client.run(getQueryPlus(), responseContext).toList()
+    );
+    Assertions.assertTrue(
+        e.getMessage().contains("failed with exception msg [connection reset 
by peer]"),
+        e.getMessage()
+    );
+    Assertions.assertTrue(e.getMessage().contains(hostName), e.getMessage());
+  }
+
+  /**
+   * An {@link HttpClient} that completes {@link 
HttpResponseHandler#handleResponse} normally, so the caller is handed
+   * a stream, and then reports a transport failure through {@link 
HttpResponseHandler#exceptionCaught} the way
+   * NettyHttpClient does when a connection drops mid-response.
+   */
+  private static class TransportFailureHttpClient implements HttpClient
+  {
+    private final RuntimeException transportFailure;
+
+    TransportFailureHttpClient(RuntimeException transportFailure)
+    {
+      this.transportFailure = transportFailure;
+    }
+
+    @Override
+    public <Intermediate, Final> ListenableFuture<Final> go(Request request, 
HttpResponseHandler<Intermediate, Final> handler)
+    {
+      return go(request, handler, null);
+    }
+
+    @Override
+    public <Intermediate, Final> ListenableFuture<Final> go(
+        Request request,
+        HttpResponseHandler<Intermediate, Final> handler,
+        Duration readTimeout
+    )
+    {
+      final HttpResponse response = new 
DefaultHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK);
+      response.headers().set(HttpHeaders.Names.CONTENT_TYPE, 
"application/json");
+      final ClientResponse<Intermediate> clientResponse =
+          handler.handleResponse(response, TestHttpClient.NOOP_TRAFFIC_COP);
+      handler.exceptionCaught(clientResponse, transportFailure);
+      @SuppressWarnings("unchecked")
+      final Final alreadyCompletedResult = (Final) clientResponse.getObj();
+      return Futures.immediateFuture(alreadyCompletedResult);
+    }
+  }
+
+  /**
+   * An {@link HttpClient} that feeds the handler a scripted response 
synchronously: the initial {@link HttpResponse}
+   * carries only headers and every element of {@code bodies} is delivered as 
an {@link HttpContent}.
+   * <p>
+   * By default, an exception thrown out of {@code handleChunk} is simply 
allowed to propagate out of {@link #go},
+   * which is adequate for testing the classification logic itself. 
Constructing with
+   * {@code routeChunkExceptionsThroughExceptionCaught=true} instead mirrors 
what
+   * {@code NettyHttpClient#messageReceived} actually does: since {@link 
DirectDruidClient#handleResponse} always
+   * returns an already-{@link ClientResponse#isFinished() finished} response, 
the future is completed with that
+   * object up front, and a later chunk's exception can only reach the handler 
via
+   * {@link HttpResponseHandler#exceptionCaught}, not by replacing the 
future's result.
+   */
+  private static class ScriptedHttpClient implements HttpClient
+  {
+    private final HttpResponseStatus status;
+    private final String contentType;
+    private final boolean routeChunkExceptionsThroughExceptionCaught;
+    private final byte[][] bodies;
+
+    ScriptedHttpClient(HttpResponseStatus status, String contentType, 
String... bodies)
+    {
+      this(status, contentType, false, 
Arrays.stream(bodies).map(StringUtils::toUtf8).toArray(byte[][]::new));
+    }
+
+    ScriptedHttpClient(
+        HttpResponseStatus status,
+        String contentType,
+        boolean routeChunkExceptionsThroughExceptionCaught,
+        byte[]... bodies
+    )
+    {
+      this.status = status;
+      this.contentType = contentType;
+      this.routeChunkExceptionsThroughExceptionCaught = 
routeChunkExceptionsThroughExceptionCaught;
+      this.bodies = bodies;
+    }
+
+    @Override
+    public <Intermediate, Final> ListenableFuture<Final> go(Request request, 
HttpResponseHandler<Intermediate, Final> handler)
+    {
+      return go(request, handler, null);
+    }
+
+    @Override
+    public <Intermediate, Final> ListenableFuture<Final> go(
+        Request request,
+        HttpResponseHandler<Intermediate, Final> handler,
+        Duration readTimeout
+    )
+    {
+      final HttpResponse response = new 
DefaultHttpResponse(HttpVersion.HTTP_1_1, status);
+      if (contentType != null) {
+        response.headers().set(HttpHeaders.Names.CONTENT_TYPE, contentType);
+      }
+      ClientResponse<Intermediate> clientResponse = 
handler.handleResponse(response, TestHttpClient.NOOP_TRAFFIC_COP);
+      final boolean initialResponseFinished = clientResponse.isFinished();
+      final Intermediate initialResponseObj = clientResponse.getObj();
+      // Netty 4 carries no body on the initial HttpResponse, so every element 
of bodies is delivered as an
+      // HttpContent, including the first. chunkNum therefore starts at 0.
+      for (int i = 0; i < bodies.length; i++) {
+        final HttpContent chunk = new 
DefaultHttpContent(Unpooled.wrappedBuffer(bodies[i]));
+        if (!routeChunkExceptionsThroughExceptionCaught) {
+          clientResponse = handler.handleChunk(clientResponse, chunk, i);
+          continue;
+        }
+        try {
+          clientResponse = handler.handleChunk(clientResponse, chunk, i);
+        }
+        catch (RuntimeException e) {
+          handler.exceptionCaught(clientResponse, e);
+          if (initialResponseFinished) {
+            // retVal was already completed by handleResponse; the caller only 
learns about this exception by
+            // reading the already-vended stream, same as real NettyHttpClient.
+            @SuppressWarnings("unchecked")
+            final Final alreadyCompletedResult = (Final) initialResponseObj;
+            return Futures.immediateFuture(alreadyCompletedResult);
+          }
+          throw e;
+        }
+      }
+      if (!routeChunkExceptionsThroughExceptionCaught) {
+        return Futures.immediateFuture(handler.done(clientResponse).getObj());
+      }
+      try {
+        return Futures.immediateFuture(handler.done(clientResponse).getObj());
+      }
+      catch (RuntimeException e) {
+        // NettyHttpClient's finishRequest() sits inside the same catch as 
handleChunk, so an exception from done()
+        // also reaches the handler through exceptionCaught after the initial 
response has completed the future.
+        handler.exceptionCaught(clientResponse, e);
+        if (initialResponseFinished) {
+          @SuppressWarnings("unchecked")
+          final Final alreadyCompletedResult = (Final) initialResponseObj;
+          return Futures.immediateFuture(alreadyCompletedResult);
+        }
+        throw e;
+      }
+    }
+  }
+
   private DirectDruidClient makeDirectDruidClient(HttpClient httpClient)
   {
-    return makeDirectDruidClient(httpClient, new NoopServiceEmitter());
+    return makeDirectDruidClient(httpClient, objectMapper, new 
NoopServiceEmitter());
   }
 
   private DirectDruidClient makeDirectDruidClient(HttpClient httpClient, 
ServiceEmitter emitter)
+  {
+    return makeDirectDruidClient(httpClient, objectMapper, emitter);
+  }
+
+  private DirectDruidClient makeDirectDruidClient(HttpClient httpClient, 
ObjectMapper clientObjectMapper)
+  {
+    return makeDirectDruidClient(httpClient, clientObjectMapper, new 
NoopServiceEmitter());
+  }
+
+  private DirectDruidClient makeDirectDruidClient(HttpClient httpClient, 
ObjectMapper clientObjectMapper, ServiceEmitter emitter)
   {
     return new DirectDruidClient(
         conglomerateRule.getConglomerate(),
         QueryRunnerTestHelper.NOOP_QUERYWATCHER,
-        objectMapper,
+        clientObjectMapper,
         httpClient,
         "http",
         hostName,
diff --git 
a/server/src/test/java/org/apache/druid/client/JsonParserIteratorTest.java 
b/server/src/test/java/org/apache/druid/client/JsonParserIteratorTest.java
index 2697c67a603..af1e42908f2 100644
--- a/server/src/test/java/org/apache/druid/client/JsonParserIteratorTest.java
+++ b/server/src/test/java/org/apache/druid/client/JsonParserIteratorTest.java
@@ -143,6 +143,62 @@ public class JsonParserIteratorTest
       });
       Assertions.assertTrue(exception.getMessage().contains("ioexception 
test"));
     }
+
+    @Test
+    public void testStructuredErrorBodyReadFailureIsNormalizedToIteratorHost()
+    {
+      // A structured error body can span chunks: init() sees START_OBJECT and 
then pulls the rest of the object
+      // through readValue(), so a later-chunk failure (a 
QueryTimeoutException that DirectDruidClient rethrows as
+      // itself) surfaces from inside that read rather than as the 
deserialized value. It has to go through
+      // convertException like every other path, or it keeps the data server's 
host instead of this iterator's.
+      final QueryTimeoutException laterChunkFailure = new 
QueryTimeoutException(
+          "timed out reading the error body",
+          "data-server-01:8100"
+      );
+      final byte[] prefix = StringUtils.toUtf8("{\"error\":\"Query 
timeout\",");
+      final InputStream truncatedErrorBody = new InputStream()
+      {
+        private int pos;
+
+        @Override
+        public int read()
+        {
+          if (pos >= prefix.length) {
+            throw laterChunkFailure;
+          }
+          return prefix[pos++] & 0xff;
+        }
+
+        @Override
+        public int read(byte[] b, int off, int len)
+        {
+          if (pos >= prefix.length) {
+            throw laterChunkFailure;
+          }
+          final int n = Math.min(len, prefix.length - pos);
+          System.arraycopy(prefix, pos, b, off, n);
+          pos += n;
+          return n;
+        }
+      };
+
+      final QueryTimeoutException exception = 
Assertions.assertThrows(QueryTimeoutException.class, () -> {
+        JsonParserIterator<Object> iterator = new JsonParserIterator<>(
+            JAVA_TYPE,
+            Futures.immediateFuture(truncatedErrorBody),
+            URL,
+            null,
+            HOST,
+            OBJECT_MAPPER
+        );
+        iterator.hasNext();
+      });
+      Assertions.assertEquals(HOST, exception.getHost());
+      Assertions.assertTrue(
+          exception.getMessage().contains("timed out reading the error body"),
+          exception.getMessage()
+      );
+    }
   }
 
   @SuppressWarnings("ResultOfMethodCallIgnored")


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

Reply via email to