FrankChen021 commented on code in PR #20313:
URL: https://github.com/apache/druid/pull/20313#discussion_r4015535831
##########
multi-stage-query/src/main/java/org/apache/druid/msq/dart/worker/DartWorkerClientImpl.java:
##########
@@ -184,6 +223,39 @@ private Pair<ServiceClient, Closeable>
getClientAndLocator(final String workerId
}
}
+ private static class RequestTrackingClient implements ServiceClient
+ {
+ private final ServiceClient delegate;
+ private final Set<ListenableFuture<?>> activeRequests;
+
+ private RequestTrackingClient(
+ final ServiceClient delegate,
+ final Set<ListenableFuture<?>> activeRequests
+ )
+ {
+ this.delegate = delegate;
+ this.activeRequests = activeRequests;
+ }
+
+ @Override
+ public <IntermediateType, FinalType> ListenableFuture<FinalType>
asyncRequest(
+ final RequestBuilder requestBuilder,
+ final HttpResponseHandler<IntermediateType, FinalType> handler
+ )
+ {
+ final ListenableFuture<FinalType> future =
delegate.asyncRequest(requestBuilder, handler);
+ activeRequests.add(future);
Review Comment:
I rechecked all 8 changed files and the current head. This finding remains
unresolved: the new postWorkOrder lock covers only work orders, but
RequestTrackingClient.asyncRequest still calls delegate.asyncRequest before
activeRequests.add. close() can race in that gap for fetch, postFinish, or stop
requests, so the returned future can be added after the close snapshot and
escape cancellation. Please serialize registration with close or cancel
immediately when the closed state is observed.
<!-- mergelens:review -->
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]