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]