Merge pull request #105 from p2p-org/missing-responses
missing responses through a WS connection
This commit is contained in:
@@ -53,6 +53,7 @@ import java.net.URI
|
||||
import java.time.Duration
|
||||
import java.time.Instant
|
||||
import java.util.Base64
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
import java.util.concurrent.Executors
|
||||
import java.util.concurrent.ScheduledFuture
|
||||
import java.util.concurrent.TimeUnit
|
||||
@@ -100,10 +101,10 @@ open class WsConnectionImpl(
|
||||
private val resetBackoffExecutor = Executors.newScheduledThreadPool(2)
|
||||
private var resetBackoffTask: ScheduledFuture<Unit>? = null
|
||||
|
||||
private val messages = Sinks
|
||||
private val subscriptionResponses = Sinks
|
||||
.many()
|
||||
.multicast()
|
||||
.directBestEffort<ResponseWSParser.WsResponse>()
|
||||
.directBestEffort<JsonRpcWsMessage>()
|
||||
|
||||
private var rpcSend = Sinks
|
||||
.many()
|
||||
@@ -114,6 +115,9 @@ open class WsConnectionImpl(
|
||||
.many()
|
||||
.multicast()
|
||||
.directBestEffort<Instant>()
|
||||
|
||||
private val currentRequests = ConcurrentHashMap<Int, Sinks.One<JsonRpcResponse>>()
|
||||
|
||||
private val sendIdSeq = AtomicInteger(IDS_START)
|
||||
private val sendExecutor = Executors.newSingleThreadExecutor()
|
||||
private var keepConnection = true
|
||||
@@ -272,41 +276,48 @@ open class WsConnectionImpl(
|
||||
|
||||
fun onMessage(msg: ResponseWSParser.WsResponse): Mono<Void> {
|
||||
return Mono.fromCallable {
|
||||
val status = messages.tryEmitNext(msg)
|
||||
if (status.isFailure) {
|
||||
if (status == Sinks.EmitResult.FAIL_ZERO_SUBSCRIBER) {
|
||||
log.debug("No subscribers to WS response")
|
||||
} else {
|
||||
log.warn("Failed to proceed with a WS message: $status")
|
||||
}
|
||||
when (msg.type) {
|
||||
ResponseWSParser.Type.RPC -> onMessageRpc(msg)
|
||||
ResponseWSParser.Type.SUBSCRIPTION -> onMessageSubscription(msg)
|
||||
}
|
||||
}.then()
|
||||
}
|
||||
|
||||
open fun getRpcResponses(): Flux<JsonRpcResponse> {
|
||||
return Flux.from(messages.asFlux())
|
||||
.publishOn(Schedulers.boundedElastic())
|
||||
.filter {
|
||||
it.type == ResponseWSParser.Type.RPC
|
||||
fun onMessageRpc(msg: ResponseWSParser.WsResponse) {
|
||||
val rpcResponse = JsonRpcResponse(
|
||||
msg.value, msg.error, msg.id, null
|
||||
)
|
||||
val sender = currentRequests.remove(msg.id.asNumber().toInt())
|
||||
if (sender == null) {
|
||||
log.warn("Unknown response received for ${msg.id}")
|
||||
} else {
|
||||
try {
|
||||
val emitResult = sender.tryEmitValue(rpcResponse)
|
||||
if (emitResult.isFailure) {
|
||||
log.debug("Response is ignored $emitResult")
|
||||
}
|
||||
} catch (t: Throwable) {
|
||||
log.warn("Response is not processed ${t.message}", t)
|
||||
}
|
||||
.map { msg ->
|
||||
JsonRpcResponse(
|
||||
msg.value, msg.error, msg.id, null
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
fun onMessageSubscription(msg: ResponseWSParser.WsResponse) {
|
||||
val subscription = JsonRpcWsMessage(
|
||||
msg.value, msg.error, msg.id.asString(),
|
||||
)
|
||||
val status = subscriptionResponses.tryEmitNext(subscription)
|
||||
if (status.isFailure) {
|
||||
if (status == Sinks.EmitResult.FAIL_ZERO_SUBSCRIBER) {
|
||||
log.debug("No subscribers to WS response")
|
||||
} else {
|
||||
log.warn("Failed to proceed with a WS message: $status")
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
open fun getSubscribeResponses(): Flux<JsonRpcWsMessage> {
|
||||
return Flux.from(messages.asFlux())
|
||||
.publishOn(Schedulers.boundedElastic())
|
||||
.filter {
|
||||
it.type == ResponseWSParser.Type.SUBSCRIPTION
|
||||
}
|
||||
.map { msg ->
|
||||
JsonRpcWsMessage(
|
||||
msg.value, msg.error, msg.id.asString(),
|
||||
)
|
||||
}
|
||||
return Flux.from(subscriptionResponses.asFlux())
|
||||
}
|
||||
|
||||
open fun callRpc(originalRequest: JsonRpcRequest): Mono<JsonRpcResponse> {
|
||||
@@ -332,7 +343,9 @@ open class WsConnectionImpl(
|
||||
}
|
||||
|
||||
fun waitForResponse(request: JsonRpcRequest, originalId: Int, startTime: Long): Mono<JsonRpcResponse> {
|
||||
val expectedId = request.id.toLong()
|
||||
val internalId = request.id.toLong()
|
||||
val onResponse = Sinks.one<JsonRpcResponse>()
|
||||
currentRequests[internalId.toInt()] = onResponse
|
||||
val noResponse = JsonRpcException(
|
||||
JsonRpcResponse.Id.from(originalId),
|
||||
JsonRpcError(
|
||||
@@ -342,13 +355,6 @@ open class WsConnectionImpl(
|
||||
false
|
||||
)
|
||||
|
||||
val response = Flux.from(getRpcResponses())
|
||||
.doOnSubscribe { sendRpc(request) }
|
||||
.filter { resp -> resp.id.asNumber() == expectedId }
|
||||
.take(Defaults.timeout)
|
||||
.take(1)
|
||||
.singleOrEmpty()
|
||||
|
||||
val failOnDisconnect = Mono.from(disconnects.asFlux())
|
||||
.flatMap {
|
||||
Mono.error<JsonRpcResponse>(
|
||||
@@ -362,11 +368,16 @@ open class WsConnectionImpl(
|
||||
)
|
||||
}
|
||||
|
||||
return response.or(failOnDisconnect)
|
||||
return Mono.from(onResponse.asMono()).or(failOnDisconnect)
|
||||
.doOnSubscribe { sendRpc(request) }
|
||||
.take(Defaults.timeout)
|
||||
.doOnNext { rpcMetrics?.timer?.record(System.nanoTime() - startTime, TimeUnit.NANOSECONDS) }
|
||||
.doOnError { rpcMetrics?.fails?.increment() }
|
||||
.map { it.copyWithId(JsonRpcResponse.Id.from(originalId)) }
|
||||
.switchIfEmpty(Mono.error(noResponse))
|
||||
.switchIfEmpty(
|
||||
Mono.fromCallable { log.warn("No response for ${request.method} ${request.params}") }.then(Mono.error(noResponse))
|
||||
)
|
||||
.doFinally { currentRequests.remove(internalId.toInt()) }
|
||||
}
|
||||
|
||||
override fun close() {
|
||||
@@ -374,5 +385,6 @@ open class WsConnectionImpl(
|
||||
keepConnection = false
|
||||
connection?.dispose()
|
||||
connection = null
|
||||
currentRequests.clear()
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user