|
|
|
@ -51,6 +51,8 @@ import java.util.concurrent.atomic.AtomicInteger; |
|
|
|
@Service |
|
|
|
public class RemoteJsInvokeService extends AbstractJsInvokeService { |
|
|
|
|
|
|
|
private static final int QUEUE_TRANSFER_DELAY = 2000; |
|
|
|
|
|
|
|
@Value("${queue.js.max_eval_requests_timeout}") |
|
|
|
private long maxEvalRequestsTimeout; |
|
|
|
|
|
|
|
@ -186,7 +188,7 @@ public class RemoteJsInvokeService extends AbstractJsInvokeService { |
|
|
|
|
|
|
|
ListenableFuture<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>> future = requestTemplate.send(new TbProtoJsQueueMsg<>(UUID.randomUUID(), jsRequestWrapper)); |
|
|
|
if (maxRequestsTimeout > 0) { |
|
|
|
future = Futures.withTimeout(future, maxRequestsTimeout, TimeUnit.MILLISECONDS, timeoutExecutorService); |
|
|
|
future = Futures.withTimeout(future, maxRequestsTimeout + QUEUE_TRANSFER_DELAY, TimeUnit.MILLISECONDS, timeoutExecutorService); |
|
|
|
} |
|
|
|
queuePushedMsgs.incrementAndGet(); |
|
|
|
Futures.addCallback(future, new FutureCallback<TbProtoQueueMsg<JsInvokeProtos.RemoteJsResponse>>() { |
|
|
|
|