diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ConfiguredUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ConfiguredUpstreams.kt index 15bbf2b7..b575c479 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ConfiguredUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ConfiguredUpstreams.kt @@ -149,12 +149,9 @@ open class ConfiguredUpstreams( ) log.info("Using ALL CHAINS (gRPC) upstream, at ${endpoint.host}:${endpoint.port}") ds.start() - .flatMapMany { - it.toFlux() - } .subscribe { - log.info("Subscribed to $it through gRPC at ${endpoint.host}:${endpoint.port}") - addUpstream(it, ds.getOrCreate(it)) + log.info("Subscribed to ${it.t1} through gRPC at ${endpoint.host}:${endpoint.port}") + addUpstream(it.t1, it.t2) } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/GrpcUpstreams.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/GrpcUpstreams.kt index 8633a623..fa702a94 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/GrpcUpstreams.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/GrpcUpstreams.kt @@ -10,7 +10,11 @@ import io.grpc.netty.NettyChannelBuilder import io.netty.handler.ssl.* import org.apache.commons.lang3.StringUtils import org.slf4j.LoggerFactory +import reactor.core.publisher.Flux import reactor.core.publisher.Mono +import reactor.core.publisher.toFlux +import reactor.util.function.Tuple2 +import reactor.util.function.Tuples import java.io.File import java.util.* import java.util.concurrent.locks.ReentrantLock @@ -29,7 +33,7 @@ class GrpcUpstreams( private var known = HashMap() private val lock = ReentrantLock() - fun start(): Mono> { + fun start(): Flux> { val channel: ManagedChannelBuilder<*> = if (auth != null && StringUtils.isNotEmpty(auth.ca)) { NettyChannelBuilder.forAddress(host, port) .useTransportSecurity() @@ -45,17 +49,18 @@ class GrpcUpstreams( val loaded = client.describe(BlockchainOuterClass.DescribeRequest.newBuilder().build()) .map { value -> val chains = ArrayList() - value.chainsList.forEach { chainDetails -> + value.chainsList.filter { + Chain.byId(it.chain.number) != Chain.UNSPECIFIED + }.map { chainDetails -> val chain = Chain.byId(chainDetails.chain.number) - if (chain != Chain.UNSPECIFIED) { - getOrCreate(chain) - .init(chainDetails) - chains.add(chain) - } + val up = getOrCreate(chain) + up.init(chainDetails) + chains.add(chain) + Tuples.of(chain, up) } - chains as List - } - .doOnError { t -> + }.flatMapMany { + it.toFlux() + }.doOnError { t -> log.error("Failed to get description from $host:$port", t) } //TODO subscribe only after receiving details @@ -93,7 +98,6 @@ class GrpcUpstreams( return if (current == null) { val created = GrpcUpstream(chain, client!!, objectMapper, upstreams.targetFor(chain)) known[chain] = created - upstreams.addUpstream(chain, created) created.start() created } else { @@ -102,4 +106,7 @@ class GrpcUpstreams( } } + fun get(chain: Chain): GrpcUpstream { + return known[chain]!! + } } \ No newline at end of file