Update AsyncTdMiddleEventBusClient.java

This commit is contained in:
Andrea Cavalli 2021-01-23 11:00:14 +01:00
parent 94ad523a75
commit d093b680f8
1 changed files with 5 additions and 18 deletions

View File

@ -1,10 +1,8 @@
package it.tdlight.tdlibsession.td.middle.client;
import io.reactivex.Completable;
import io.vertx.circuitbreaker.CircuitBreakerOptions;
import io.vertx.core.eventbus.DeliveryOptions;
import io.vertx.core.json.JsonObject;
import io.vertx.reactivex.circuitbreaker.CircuitBreaker;
import io.vertx.reactivex.core.AbstractVerticle;
import io.vertx.reactivex.core.eventbus.Message;
import it.tdlight.jni.TdApi;
@ -80,17 +78,6 @@ public class AsyncTdMiddleEventBusClient extends AbstractVerticle implements Asy
@Override
public Completable rxStart() {
CircuitBreaker startBreaker = CircuitBreaker.create("bot-" + botAddress + "-server-online-check-circuit-breaker", vertx,
new CircuitBreakerOptions().setMaxFailures(1).setMaxRetries(4).setTimeout(30000)
)
.retryPolicy(policy -> 4000L)
.openHandler(closed -> {
logger.error("Circuit opened! " + botAddress);
})
.closeHandler(closed -> {
logger.error("Circuit closed! " + botAddress);
});
return Mono
.fromCallable(() -> {
var botAddress = config().getString("botAddress");
@ -111,7 +98,7 @@ public class AsyncTdMiddleEventBusClient extends AbstractVerticle implements Asy
this.initTime = System.currentTimeMillis();
return null;
})
.then(startBreaker.rxExecute(future -> {
.then(Mono.create(future -> {
logger.debug("Waiting for " + botAddress + ".readyToStart");
var readyToStartConsumer = cluster.getEventBus().<byte[]>consumer(botAddress + ".readyToStart");
readyToStartConsumer.handler((Message<byte[]> pingMsg) -> {
@ -137,7 +124,8 @@ public class AsyncTdMiddleEventBusClient extends AbstractVerticle implements Asy
.rxRequest(botAddress + ".isWorking", EMPTY, deliveryOptionsWithTimeout)
.as(MonoUtils::toMono)
.doOnError(ex -> logger.error("Failed to request isWorking", ex)))
.subscribe(v -> {}, future::fail, future::complete);
.subscribeOn(Schedulers.single())
.subscribe(v -> {}, future::error, future::success);
});
readyToStartConsumer
@ -145,9 +133,8 @@ public class AsyncTdMiddleEventBusClient extends AbstractVerticle implements Asy
.as(MonoUtils::toMono)
.doOnSuccess(s -> logger.debug("Requesting tdlib.remoteclient.clients.deploy.existing"))
.then(Mono.fromCallable(() -> cluster.getEventBus().publish("tdlib.remoteclient.clients.deploy", botAddress, deliveryOptions)))
.subscribeOn(Schedulers.single())
.subscribe(v -> {}, future::fail);
}).as(MonoUtils::toMono)
.subscribe(v -> {}, future::error);
})
)
.onErrorMap(ex -> {
logger.error("Failure when starting bot " + botAddress, ex);