diff --git a/demo/test.html b/demo/test.html
deleted file mode 100644
index 9a44c15a..00000000
--- a/demo/test.html
+++ /dev/null
@@ -1,13 +0,0 @@
-
-
-
-
-
-
\ No newline at end of file
diff --git a/docs/reference-configuration.adoc b/docs/reference-configuration.adoc
index d71e8589..0d691617 100644
--- a/docs/reference-configuration.adoc
+++ b/docs/reference-configuration.adoc
@@ -678,8 +678,7 @@ configuration, and may be omitted for most of the situations.
| `role`
| no
-| `primary` (default), `secondary` or `fallback`.
-First it makes the requests to the upstreams with role `primary`, then if none are available to upstreams with role `secondary`.
+| `standard` (default) or `fallback`.
Fallback role mean that the upstream is used only after other upstreams failed or didn't return quorum
| `chain`
diff --git a/src/main/kotlin/io/emeraldpay/dshackle/proxy/WebsocketHandler.kt b/src/main/kotlin/io/emeraldpay/dshackle/proxy/WebsocketHandler.kt
index b1fb0d7a..3882d05b 100644
--- a/src/main/kotlin/io/emeraldpay/dshackle/proxy/WebsocketHandler.kt
+++ b/src/main/kotlin/io/emeraldpay/dshackle/proxy/WebsocketHandler.kt
@@ -27,6 +27,7 @@ import io.emeraldpay.dshackle.rpc.NativeSubscribe
import io.emeraldpay.etherjar.rpc.json.RequestJson
import io.emeraldpay.etherjar.rpc.json.ResponseJson
import io.netty.buffer.ByteBufInputStream
+import io.netty.buffer.Unpooled
import org.reactivestreams.Publisher
import org.slf4j.LoggerFactory
import reactor.core.publisher.Flux
@@ -73,8 +74,9 @@ class WebsocketHandler(
val eventHandler = accessHandler.start(req, routeConfig.blockchain)
val responses = respond(routeConfig.blockchain, control, requests, eventHandler)
+ .map { Unpooled.wrappedBuffer(it.toByteArray()) }
- resp.sendString(responses, Charsets.UTF_8)
+ resp.send(responses)
.then()
}
}
diff --git a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt
index 912f7c64..73d52cd4 100644
--- a/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt
+++ b/src/main/kotlin/io/emeraldpay/dshackle/rpc/NativeCall.kt
@@ -284,7 +284,11 @@ open class NativeCall(
Mono.just(ctx).flatMap(this::executeOnRemote)
)
.onErrorResume {
- Mono.just(CallResult.fail(ctx.id, ctx.nonce, it))
+ if (it is CallFailure) {
+ Mono.just(CallResult.fail(it.id, ctx.nonce, it.reason))
+ } else {
+ Mono.just(CallResult.fail(ctx.id, ctx.nonce, it))
+ }
}
}
@@ -308,14 +312,21 @@ open class NativeCall(
CallResult.ok(ctx.id, ctx.nonce, bytes, it.signature, ctx.upstream.getId())
}
.onErrorResume { t ->
- Mono.just(CallResult.fail(ctx.id, ctx.nonce, t))
+ val failure = when (t) {
+ is CallFailure -> CallResult.fail(t.id, ctx.nonce, t.reason)
+ is JsonRpcException -> CallResult.fail(ctx.id, ctx.nonce, t.error.code, t.error.message)
+ else -> CallResult.fail(ctx.id, ctx.nonce, t)
+ }
+ Mono.just(failure)
}
.switchIfEmpty(
Mono.fromSupplier {
counter.get().let { attempts ->
CallResult.fail(
- ctx.id, ctx.nonce,
- CallError(1, "No response or no available upstream for ${ctx.payload.method}", null)
+ ctx.id,
+ ctx.nonce,
+ 1,
+ errorMessage(attempts, ctx.payload.method)
).also {
countFailure(attempts, ctx)
}
@@ -480,15 +491,7 @@ open class NativeCall(
is JsonRpcException -> CallError(t.id.asNumber().toInt(), t.error.message, t.error)
is RpcException -> CallError(t.code, t.rpcMessage, null)
is CallFailure -> CallError(t.id, t.reason.message ?: "Upstream Error", null)
- else -> {
- // May only happen if it's an unhandled exception.
- // In this case try to find a meaningless details in the stack. Most important reason for doing that is to find an ID of the request
- if (t.cause != null) {
- from(t.cause!!)
- } else {
- CallError(1, t.message ?: "Upstream Error", null)
- }
- }
+ else -> CallError(1, t.message ?: "Upstream Error", null)
}
}
}
@@ -507,8 +510,8 @@ open class NativeCall(
return CallResult(id, nonce, result, null, signature, upstreamId)
}
- fun fail(id: Int, nonce: Long?, error: CallError): CallResult {
- return CallResult(id, nonce, null, error, null, null)
+ fun fail(id: Int, nonce: Long?, errorCore: Int, errorMessage: String): CallResult {
+ return CallResult(id, nonce, null, CallError(errorCore, errorMessage, null), null, null)
}
fun fail(id: Int, nonce: Long?, error: Throwable): CallResult {
diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/DefaultUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/DefaultUpstream.kt
index 0f967fba..451a53b4 100644
--- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/DefaultUpstream.kt
+++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/DefaultUpstream.kt
@@ -98,11 +98,6 @@ abstract class DefaultUpstream(
}
fun statusByLag(lag: Long, proposed: UpstreamAvailability): UpstreamAvailability {
- if (options.disableValidation == true) {
- // if we specifically told that this upstream should be _always valid_ then skip
- // the status calculation and trust the proposed value as is
- return proposed
- }
return if (proposed == UpstreamAvailability.OK) {
when {
lag > 6 -> UpstreamAvailability.SYNCING
diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnection.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnection.kt
index 5373e796..c30f3f89 100644
--- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnection.kt
+++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/WsConnection.kt
@@ -405,8 +405,7 @@ open class WsConnection(
)
val response = Flux.from(rpcReceive.asFlux())
- // send the request _after_ WS subscribes to the responses, otherwise the response may come before the actual subscription and be lost
- .doOnRequest { sendRpc(request) }
+ .doOnSubscribe { sendRpc(request) }
.filter { resp -> resp.id.asNumber() == expectedId }
.take(Defaults.timeout)
.take(1)
diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcGrpcClient.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcGrpcClient.kt
index b525b1c2..8acff9e2 100644
--- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcGrpcClient.kt
+++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcGrpcClient.kt
@@ -25,8 +25,6 @@ import io.emeraldpay.dshackle.reader.Reader
import io.emeraldpay.dshackle.upstream.signature.ResponseSigner
import io.emeraldpay.etherjar.rpc.RpcException
import io.emeraldpay.etherjar.rpc.RpcResponseError
-import io.grpc.StatusRuntimeException
-import org.apache.commons.lang3.time.StopWatch
import org.slf4j.LoggerFactory
import reactor.core.publisher.Mono
import java.util.concurrent.TimeUnit
@@ -34,7 +32,7 @@ import java.util.concurrent.TimeUnit
class JsonRpcGrpcClient(
private val stub: ReactorBlockchainGrpc.ReactorBlockchainStub,
private val chain: Chain,
- private val metrics: RpcMetrics?,
+ private val metrics: RpcMetrics
) {
companion object {
@@ -48,11 +46,11 @@ class JsonRpcGrpcClient(
class Executor(
private val stub: ReactorBlockchainGrpc.ReactorBlockchainStub,
private val chain: Chain,
- private val metrics: RpcMetrics?
+ private val metrics: RpcMetrics
) : Reader {
override fun read(key: JsonRpcRequest): Mono {
- val timer = StopWatch()
+ var startTime: Long = 0
val req = BlockchainOuterClass.NativeCallRequest.newBuilder()
.setChainValue(chain.id)
@@ -68,58 +66,39 @@ class JsonRpcGrpcClient(
req.addItems(reqItem.build())
return Mono.just(key)
- .doOnNext { timer.start() }
- .flatMap {
+ .doOnNext {
+ startTime = System.nanoTime()
+ }.flatMap {
stub.nativeCall(req.build())
.single()
- .onErrorResume(::handleError)
- .flatMap(::handleResponse)
+ .flatMap { resp ->
+ if (resp.succeed) {
+ val bytes = resp.payload.toByteArray()
+ val signature = if (resp.hasSignature()) {
+ extractSignature(resp.signature)
+ } else {
+ null
+ }
+ Mono.just(JsonRpcResponse(bytes, null, JsonRpcResponse.NumberId(0), signature))
+ } else {
+ metrics.fails.increment()
+ Mono.error(
+ RpcException(
+ RpcResponseError.CODE_UPSTREAM_CONNECTION_ERROR,
+ resp.errorMessage
+ )
+ )
+ }
+ }
}
.doOnNext {
- if (timer.isStarted) {
- metrics?.timer?.record(timer.getTime(TimeUnit.NANOSECONDS), TimeUnit.NANOSECONDS)
+ if (startTime > 0) {
+ val now = System.nanoTime()
+ metrics.timer.record(now - startTime, TimeUnit.NANOSECONDS)
}
}
}
- fun handleResponse(resp: BlockchainOuterClass.NativeCallReplyItem): Mono =
- if (resp.succeed) {
- val bytes = resp.payload.toByteArray()
- val signature = if (resp.hasSignature()) {
- extractSignature(resp.signature)
- } else {
- null
- }
- Mono.just(JsonRpcResponse(bytes, null, JsonRpcResponse.NumberId(0), signature))
- } else {
- metrics?.fails?.increment()
- Mono.error(
- RpcException(
- RpcResponseError.CODE_UPSTREAM_CONNECTION_ERROR,
- resp.errorMessage
- )
- )
- }
-
- fun handleError(t: Throwable): Mono {
- metrics?.fails?.increment()
- return when (t) {
- is StatusRuntimeException -> Mono.error(
- RpcException(
- RpcResponseError.CODE_UPSTREAM_CONNECTION_ERROR,
- "Remote status code: ${t.status.code.name}"
- )
- )
-
- else -> Mono.error(
- RpcException(
- RpcResponseError.CODE_UPSTREAM_CONNECTION_ERROR,
- "Other connection error"
- )
- )
- }
- }
-
fun extractSignature(resp: NativeCallReplySignature?): ResponseSigner.Signature? {
if (resp == null || resp.signature == null || resp.signature.isEmpty || resp.upstreamId == null || resp.upstreamId.isEmpty()) {
return null
diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClient.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClient.kt
index 154ea0c4..8535bdc2 100644
--- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClient.kt
+++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClient.kt
@@ -24,13 +24,10 @@ import io.netty.handler.codec.http.HttpHeaderNames
import io.netty.handler.codec.http.HttpHeaders
import io.netty.handler.ssl.SslContextBuilder
import io.netty.resolver.DefaultAddressResolverGroup
-import org.apache.commons.lang3.time.StopWatch
import org.slf4j.LoggerFactory
import reactor.core.publisher.Mono
import reactor.netty.http.client.HttpClient
import reactor.netty.resources.ConnectionProvider
-import reactor.util.function.Tuple2
-import reactor.util.function.Tuples
import java.io.ByteArrayInputStream
import java.security.KeyStore
import java.security.cert.CertificateFactory
@@ -38,7 +35,6 @@ import java.security.cert.X509Certificate
import java.util.Base64
import java.util.concurrent.TimeUnit
import java.util.function.Consumer
-import java.util.function.Function
/**
* JSON RPC client
@@ -93,95 +89,52 @@ class JsonRpcHttpClient(
this.httpClient = build
}
- fun execute(request: ByteArray): Mono> {
+ fun execute(request: ByteArray): Mono {
val response = httpClient
.post()
.uri(target)
.send(Mono.just(request).map { Unpooled.wrappedBuffer(it) })
return response.response { header, bytes ->
- val statusCode = header.status().code()
- bytes.aggregate().asByteArray().map {
- Tuples.of(statusCode, it)
+ if (header.status().code() != 200) {
+ Mono.error(
+ JsonRpcException(
+ JsonRpcResponse.NumberId(-2),
+ JsonRpcError(
+ RpcResponseError.CODE_UPSTREAM_INVALID_RESPONSE,
+ "HTTP Code: ${header.status().code()}"
+ )
+ )
+ )
+ } else {
+ bytes.aggregate().asByteArray()
}
}.single()
}
override fun read(key: JsonRpcRequest): Mono {
- val startTime = StopWatch()
+ var startTime: Long = 0
return Mono.just(key)
.map(JsonRpcRequest::toJson)
- .doOnNext { startTime.start() }
+ .doOnNext {
+ startTime = System.nanoTime()
+ }
.flatMap(this@JsonRpcHttpClient::execute)
.doOnNext {
- if (startTime.isStarted) {
- metrics.timer.record(startTime.nanoTime, TimeUnit.NANOSECONDS)
+ if (startTime > 0) {
+ val now = System.nanoTime()
+ metrics.timer.record(now - startTime, TimeUnit.NANOSECONDS)
}
}
- .transform(asJsonRpcResponse(key))
- .transform(convertErrors(key))
- .transform(throwIfError())
- }
-
- /**
- * The subscribers expect to catch an exception if the response contains JSON RPC Error. Convert it here to JsonRpcException
- */
- private fun throwIfError(): Function, Mono> {
- return Function { resp ->
- resp.flatMap {
- if (it.hasError()) {
- Mono.error(JsonRpcException(it.id, it.error!!))
- } else {
- Mono.just(it)
- }
- }
- }
- }
-
- /**
- * Convert internal exceptions to standard JsonRpcException
- */
- private fun convertErrors(key: JsonRpcRequest): Function, Mono> {
- return Function { resp ->
- resp.onErrorResume { t ->
+ .map(parser::parse)
+ .onErrorResume { t ->
val err = when (t) {
- is RpcException -> JsonRpcException.from(t)
- is JsonRpcException -> t
- else -> JsonRpcException(key.id, t.message ?: t.javaClass.name)
+ is RpcException -> JsonRpcResponse.error(t.code, t.rpcMessage)
+ is JsonRpcException -> JsonRpcResponse.error(t.error, JsonRpcResponse.NumberId(1))
+ else -> JsonRpcResponse.error(1, t.message ?: t.javaClass.name)
}
- // here we're measure the internal errors, not upstream errors
metrics.fails.increment()
- Mono.error(err)
+ Mono.just(err)
}
- }
- }
-
- /**
- * Process response from the upstream and convert it to JsonRpcResponse.
- * The input is a pair of (Http Status Code, Http Response Body)
- */
- private fun asJsonRpcResponse(key: JsonRpcRequest): Function>, Mono> {
- return Function { resp ->
- resp.map {
- val parsed = parser.parse(it.t2)
- val statusCode = it.t1
- if (statusCode != 200) {
- if (parsed.hasError() && parsed.error!!.code != RpcResponseError.CODE_UPSTREAM_INVALID_RESPONSE) {
- // extracted the error details from the HTTP Body
- parsed
- } else {
- // here we got a valid response with ERROR as HTTP Status Code. We assume that HTTP Status has
- // a higher priority so return an error here anyway
- JsonRpcResponse.error(
- RpcResponseError.CODE_UPSTREAM_INVALID_RESPONSE,
- "HTTP Code: $statusCode",
- JsonRpcResponse.NumberId(key.id)
- )
- }
- } else {
- parsed
- }
- }
- }
}
}
diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/ResponseParser.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/ResponseParser.kt
index 2e4927b3..56f106a0 100644
--- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/ResponseParser.kt
+++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/ResponseParser.kt
@@ -21,7 +21,6 @@ import com.fasterxml.jackson.core.JsonParser
import com.fasterxml.jackson.core.JsonToken
import io.emeraldpay.dshackle.Global
import io.emeraldpay.etherjar.rpc.RpcResponseError
-import org.apache.commons.lang3.StringUtils
import org.slf4j.LoggerFactory
import java.io.IOException
@@ -59,18 +58,15 @@ abstract class ResponseParser {
} catch (e: JsonParseException) {
log.warn("Failed to parse JSON from upstream: ${e.message}")
}
- return if (state.isReady) {
- state
- } else {
- log.debug("Failed to parse `${StringUtils.abbreviateMiddle(String(json), "...", 200)}` JSON")
- state.copy(
- result = null,
- error = JsonRpcError(
- RpcResponseError.CODE_UPSTREAM_INVALID_RESPONSE,
- "Invalid JSON structure: never finalized"
- )
- )
+ if (state.isReady) {
+ return state
}
+ return Preparsed(
+ error = JsonRpcError(
+ RpcResponseError.CODE_UPSTREAM_INVALID_RESPONSE,
+ "Invalid JSON structure: never finalized"
+ )
+ )
}
open fun process(parser: JsonParser, json: ByteArray, field: String, state: Preparsed): Preparsed {
diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParser.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParser.kt
index 02a7416b..87c76746 100644
--- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParser.kt
+++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParser.kt
@@ -44,16 +44,6 @@ class ResponseWSParser : ResponseParser() {
state.error
)
}
- if (state.error != null) {
- return WsResponse(
- // we don't have any real option because it's just an invalid value and can be anything,
- // so let's suppose its Type as RPC as a most likely scenario
- Type.RPC,
- state.id ?: JsonRpcResponse.Id.from(0),
- null,
- state.error
- )
- }
throw IllegalStateException("State is not ready")
}
diff --git a/src/test/groovy/io/emeraldpay/dshackle/config/CacheConfigReaderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/config/CacheConfigReaderSpec.groovy
index 99305cf9..8b0901ed 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/config/CacheConfigReaderSpec.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/config/CacheConfigReaderSpec.groovy
@@ -23,7 +23,7 @@ class CacheConfigReaderSpec extends Specification {
def "Read full"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/cache-redis-full.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("cache-redis-full.yaml")
when:
def act = reader.read(config)
@@ -39,7 +39,7 @@ class CacheConfigReaderSpec extends Specification {
def "Read disabled"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/cache-redis-disabled.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("cache-redis-disabled.yaml")
when:
def act = reader.read(config)
diff --git a/src/test/groovy/io/emeraldpay/dshackle/config/HealthConfigReaderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/config/HealthConfigReaderSpec.groovy
index 9006581b..a340c75c 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/config/HealthConfigReaderSpec.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/config/HealthConfigReaderSpec.groovy
@@ -24,7 +24,7 @@ class HealthConfigReaderSpec extends Specification {
def "Read empty"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/dshackle-health-empty.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("dshackle-health-empty.yaml")
when:
def act = reader.read(config)
@@ -38,7 +38,7 @@ class HealthConfigReaderSpec extends Specification {
def "Read single"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/dshackle-health-1.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("dshackle-health-1.yaml")
when:
def act = reader.read(config)
@@ -58,7 +58,7 @@ class HealthConfigReaderSpec extends Specification {
def "Read multiple"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/dshackle-health-2.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("dshackle-health-2.yaml")
when:
def act = reader.read(config)
diff --git a/src/test/groovy/io/emeraldpay/dshackle/config/MainConfigReaderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/config/MainConfigReaderSpec.groovy
index 8a0c0eab..89d48c42 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/config/MainConfigReaderSpec.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/config/MainConfigReaderSpec.groovy
@@ -22,12 +22,11 @@ import spock.lang.Specification
class MainConfigReaderSpec extends Specification {
- def "Read full config with inclusion"() {
+ MainConfigReader reader = new MainConfigReader(TestingCommons.fileResolver())
+
+ def "Read full config"() {
setup:
- // the File Resolver should be able to resolve/include files
- MainConfigReader reader = new MainConfigReader(TestingCommons.fileResolver())
- // note that it references another config to be included
- def config = this.class.getClassLoader().getResourceAsStream("configs/dshackle-full.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("dshackle-full.yaml")
when:
def act = reader.read(config)
diff --git a/src/test/groovy/io/emeraldpay/dshackle/config/MonitoringConfigReaderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/config/MonitoringConfigReaderSpec.groovy
index eee026ca..a1fa8f67 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/config/MonitoringConfigReaderSpec.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/config/MonitoringConfigReaderSpec.groovy
@@ -23,7 +23,7 @@ class MonitoringConfigReaderSpec extends Specification {
def "Read basic monitoring config"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/dshackle-monitoring-basic.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("dshackle-monitoring-basic.yaml")
when:
def act = reader.read(config)
diff --git a/src/test/groovy/io/emeraldpay/dshackle/config/ProxyConfigReaderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/config/ProxyConfigReaderSpec.groovy
index 473247df..8d2f827d 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/config/ProxyConfigReaderSpec.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/config/ProxyConfigReaderSpec.groovy
@@ -24,7 +24,7 @@ class ProxyConfigReaderSpec extends Specification {
def "Read basic proxy config"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/dshackle-proxy-basic.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("dshackle-proxy-basic.yaml")
when:
def act = reader.read(config)
@@ -42,7 +42,7 @@ class ProxyConfigReaderSpec extends Specification {
def "Read proxy config with websocket disabled"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/dshackle-proxy-no-ws.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("dshackle-proxy-no-ws.yaml")
when:
def act = reader.read(config)
@@ -53,7 +53,7 @@ class ProxyConfigReaderSpec extends Specification {
def "Read proxy config with two elements"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/dshackle-proxy-two.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("dshackle-proxy-two.yaml")
when:
def act = reader.read(config)
@@ -73,7 +73,7 @@ class ProxyConfigReaderSpec extends Specification {
def "Read max proxy config"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/dshackle-proxy-max.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("dshackle-proxy-max.yaml")
when:
def act = reader.read(config)
diff --git a/src/test/groovy/io/emeraldpay/dshackle/config/UpstreamsConfigReaderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/config/UpstreamsConfigReaderSpec.groovy
index d115c12e..5b4589a8 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/config/UpstreamsConfigReaderSpec.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/config/UpstreamsConfigReaderSpec.groovy
@@ -25,7 +25,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse standard config"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-basic.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-basic.yaml")
when:
def act = reader.read(config)
then:
@@ -73,7 +73,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse websocket-only config"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-ws-only.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-ws-only.yaml")
when:
def act = reader.read(config)
then:
@@ -98,7 +98,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse full defined websocket config"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-ws-full.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-ws-full.yaml")
when:
def act = reader.read(config)
then:
@@ -125,7 +125,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse bitcoin upstreams"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-bitcoin.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-bitcoin.yaml")
when:
def act = reader.read(config)
then:
@@ -173,7 +173,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse bitcoin upstreams with esplora"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-bitcoin-esplora.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-bitcoin-esplora.yaml")
when:
def act = reader.read(config)
then:
@@ -202,7 +202,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse ds config"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-ds.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-ds.yaml")
when:
def act = reader.read(config)
then:
@@ -225,7 +225,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse config with labels"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-labels.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-labels.yaml")
when:
def act = reader.read(config)
then:
@@ -246,7 +246,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse config with options"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-options.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-options.yaml")
when:
def act = reader.read(config)
then:
@@ -262,7 +262,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse config without defaults"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-no-defaults.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-no-defaults.yaml")
when:
def act = reader.read(config)
then:
@@ -285,7 +285,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse config with methods"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-methods.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-methods.yaml")
when:
def act = reader.read(config)
then:
@@ -305,7 +305,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse config with methods and quorum"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-methods-quorum.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-methods-quorum.yaml")
when:
def act = reader.read(config)
then:
@@ -328,7 +328,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse config with invalid ids"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-no-id.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-no-id.yaml")
when:
def act = reader.read(config)
then:
@@ -355,7 +355,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse config without fallback role"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-basic.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-basic.yaml")
when:
def act = reader.read(config)
then:
@@ -367,7 +367,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse config with fallback role"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-roles.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-roles.yaml")
when:
def act = reader.read(config)
then:
@@ -379,7 +379,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse config with secondary role"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-roles-2.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-roles-2.yaml")
when:
def act = reader.read(config)
then:
@@ -392,7 +392,7 @@ class UpstreamsConfigReaderSpec extends Specification {
def "Parse config with invalid role"() {
setup:
- def config = this.class.getClassLoader().getResourceAsStream("configs/upstreams-roles-invalid.yaml")
+ def config = this.class.getClassLoader().getResourceAsStream("upstreams-roles-invalid.yaml")
when:
def act = reader.read(config)
then:
diff --git a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy
index 8a45dc30..0f039786 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/rpc/NativeCallSpec.groovy
@@ -34,8 +34,6 @@ import io.emeraldpay.dshackle.upstream.Selector
import io.emeraldpay.dshackle.upstream.MultistreamHolder
import io.emeraldpay.dshackle.upstream.calls.DefaultEthereumMethods
import io.emeraldpay.dshackle.upstream.calls.ManagedCallMethods
-import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcError
-import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcException
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcRequest
import io.emeraldpay.dshackle.upstream.rpcclient.JsonRpcResponse
import io.emeraldpay.dshackle.upstream.signature.ResponseSigner
@@ -155,35 +153,6 @@ class NativeCallSpec extends Specification {
.verify(Duration.ofSeconds(1))
}
- def "Returns error details from remote"() {
- setup:
- def quorum = new AlwaysQuorum()
-
- def nativeCall = nativeCall()
- nativeCall.quorumReaderFactory = Mock(QuorumReaderFactory) {
- 1 * create(_, _, _) >> Mock(Reader) {
- 1 * read(new JsonRpcRequest("eth_test", [], 10)) >> Mono.error(
- new JsonRpcException(JsonRpcResponse.Id.from(12), new JsonRpcError(-32123, "Foo Bar", "Foo Bar Baz"))
- )
- }
- }
- def call = new NativeCall.ValidCallContext(12, 10, TestingCommons.multistream(TestingCommons.api()), Selector.empty, quorum,
- new NativeCall.ParsedCallDetails("eth_test", []))
-
- when:
- def resp = nativeCall.executeOnRemote(call).block(Duration.ofSeconds(1))
- then:
- resp.isError()
- with(resp.getError()) {
- message == "Foo Bar"
- upstreamError != null
- with (upstreamError) {
- code == -32123
- details == "Foo Bar Baz"
- }
- }
- }
-
def "Packs call exception into response with id"() {
setup:
def nativeCall = nativeCall()
diff --git a/src/test/groovy/io/emeraldpay/dshackle/test/MockGrpcServer.groovy b/src/test/groovy/io/emeraldpay/dshackle/test/MockGrpcServer.groovy
index f1d3689f..3e046e9b 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/test/MockGrpcServer.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/test/MockGrpcServer.groovy
@@ -26,6 +26,14 @@ class MockGrpcServer {
GrpcCleanupRule grpcCleanup = new GrpcCleanupRule()
+ ReactorBlockchainGrpc.ReactorBlockchainStub clientForServer(ReactorBlockchainGrpc.BlockchainImplBase impl) {
+ String serverName = InProcessServerBuilder.generateName()
+ grpcCleanup.register(InProcessServerBuilder
+ .forName(serverName).directExecutor().addService(impl).build().start());
+ def channel = grpcCleanup.register(InProcessChannelBuilder.forName(serverName).directExecutor().build())
+ return ReactorBlockchainGrpc.newReactorStub(channel)
+ }
+
ReactorBlockchainGrpc.ReactorBlockchainStub clientForServer(BlockchainGrpc.BlockchainImplBase impl){
String serverName = InProcessServerBuilder.generateName()
grpcCleanup.register(InProcessServerBuilder
diff --git a/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy b/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy
index a7405aa4..a6e4c114 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/test/TestingCommons.groovy
@@ -108,7 +108,7 @@ class TestingCommons {
}
static FileResolver fileResolver() {
- return new FileResolver(new File("src/test/resources/configs"))
+ return new FileResolver(new File("src/test/resources"))
}
static BlockContainer blockForEthereum(Long height) {
diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionSpec.groovy
index 54849052..18350845 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionSpec.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/ethereum/WsConnectionSpec.groovy
@@ -92,7 +92,7 @@ class WsConnectionSpec extends Specification {
it.id.asNumber() == 15L && Global.objectMapper.readValue(it.result, TransactionJson) == tx
}
.expectComplete()
- .verify(Duration.ofSeconds(1))
+ .verify(Duration.ofSeconds(5))
}
def "Makes a RPC call - return null"() {
@@ -115,7 +115,7 @@ class WsConnectionSpec extends Specification {
it.resultAsRawString == 'null'
}
.expectComplete()
- .verify(Duration.ofSeconds(1))
+ .verify(Duration.ofSeconds(5))
}
def "Makes a RPC call - return error"() {
@@ -140,6 +140,6 @@ class WsConnectionSpec extends Specification {
it.error.code == RpcResponseError.CODE_METHOD_NOT_EXIST && it.error.message == "test"
}
.expectComplete()
- .verify(Duration.ofSeconds(1))
+ .verify(Duration.ofSeconds(5))
}
}
diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcGrpcClientSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcGrpcClientSpec.groovy
deleted file mode 100644
index 318a436f..00000000
--- a/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcGrpcClientSpec.groovy
+++ /dev/null
@@ -1,88 +0,0 @@
-package io.emeraldpay.dshackle.upstream.rpcclient
-
-import com.google.protobuf.ByteString
-import io.emeraldpay.api.proto.BlockchainGrpc
-import io.emeraldpay.api.proto.BlockchainOuterClass
-import io.emeraldpay.api.proto.Common
-import io.emeraldpay.dshackle.test.MockGrpcServer
-import io.emeraldpay.dshackle.upstream.Selector
-import io.emeraldpay.etherjar.rpc.RpcException
-import io.emeraldpay.etherjar.rpc.RpcResponseError
-import io.emeraldpay.dshackle.Chain
-import io.grpc.stub.StreamObserver
-import spock.lang.Specification
-
-import java.time.Duration
-import java.util.concurrent.atomic.AtomicReference
-
-class JsonRpcGrpcClientSpec extends Specification {
-
- def "Makes a request"() {
- setup:
- def mockGrpc = new MockGrpcServer()
- def requested = new AtomicReference()
-
- def grpc = mockGrpc.clientForServer(new BlockchainGrpc.BlockchainImplBase() {
- @Override
- void nativeCall(BlockchainOuterClass.NativeCallRequest request, StreamObserver responseObserver) {
- requested.set(request)
- responseObserver.onNext(BlockchainOuterClass.NativeCallReplyItem.newBuilder()
- .setId(1)
- .setSucceed(true)
- .setPayload(ByteString.copyFromUtf8("\"hello world!\""))
- .build())
- responseObserver.onCompleted()
- }
- })
- def client = new JsonRpcGrpcClient(
- grpc, Chain.BITCOIN, null
- ).getReader()
-
- when:
- def act = client.read(
- new JsonRpcRequest("test", [])
- ).block(Duration.ofSeconds(1))
-
- then:
- !act.hasError()
- act.resultAsProcessedString == "hello world!"
-
- requested.get() == BlockchainOuterClass.NativeCallRequest.newBuilder()
- .setChain(Common.ChainRef.CHAIN_BITCOIN)
- .addAllItems([
- BlockchainOuterClass.NativeCallItem.newBuilder()
- .setId(1)
- .setMethod("test")
- .setPayload(ByteString.copyFromUtf8("[]"))
- .build()
- ])
- .build()
- }
-
- def "Return error on HTTP error"() {
- setup:
- def mockGrpc = new MockGrpcServer()
- def grpc = mockGrpc.clientForServer(new BlockchainGrpc.BlockchainImplBase() {
- @Override
- void nativeCall(BlockchainOuterClass.NativeCallRequest request, StreamObserver responseObserver) {
- responseObserver.onError(new IllegalStateException("fail"))
- }
- })
- def client = new JsonRpcGrpcClient(
- grpc, Chain.BITCOIN, null
- ).getReader()
-
- when:
- client.read(
- new JsonRpcRequest("test", [])
- ).block(Duration.ofSeconds(1))
-
- then:
- def t = thrown(RpcException)
- with(t.error) {
- message == "Remote status code: UNKNOWN"
- it.code == RpcResponseError.CODE_UPSTREAM_CONNECTION_ERROR
- }
-
- }
-}
diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClientSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClientSpec.groovy
index 549e81fe..1e02cdc1 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClientSpec.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/JsonRpcHttpClientSpec.groovy
@@ -26,7 +26,6 @@ import org.mockserver.model.HttpRequest
import org.mockserver.model.HttpResponse
import org.mockserver.model.MediaType
import org.springframework.util.SocketUtils
-import reactor.core.Exceptions
import reactor.test.StepVerifier
import spock.lang.Specification
@@ -85,7 +84,7 @@ class JsonRpcHttpClientSpec extends Specification {
.withBody("pong")
)
when:
- def act = client.execute("ping".bytes).map { new String(it.t2) }
+ def act = client.execute("ping".bytes).map { new String(it) }
then:
StepVerifier.create(act)
.expectNext("pong")
@@ -112,44 +111,13 @@ class JsonRpcHttpClientSpec extends Specification {
.withBody("pong")
)
when:
- def act = client.read(
- new JsonRpcRequest("ping", [])
- ).block(Duration.ofSeconds(1))
+ def act = client.execute("ping".bytes).map { new String(it) }
then:
- def t = thrown(RuntimeException) // reactor.core.Exceptions$ReactiveException
- t.cause instanceof JsonRpcException
- with(((JsonRpcException)t.cause).error) {
- code == RpcResponseError.CODE_UPSTREAM_INVALID_RESPONSE
- message == "HTTP Code: 500"
- }
- }
-
- def "Tries to extract message if HTTP error if it still contains a JSON RPC message"() {
- setup:
- def client = new JsonRpcHttpClient("localhost:${port}", metrics, null, null)
-
- mockServer.when(
- HttpRequest.request()
- ).respond(
- HttpResponse.response()
- .withStatusCode(500)
- .withBody('{' +
- '"jsonrpc": "2.0", ' +
- '"id": 1, ' +
- '"error": {"code": -32603, "message": "Something happened"}' +
- '}')
- )
- when:
- def act = client.read(
- new JsonRpcRequest("ping", [])
- ).block(Duration.ofSeconds(1))
- then:
- def t = thrown(RuntimeException) // reactor.core.Exceptions$ReactiveException
- t.cause instanceof JsonRpcException
- with(((JsonRpcException)t.cause).error) {
- code == -32603
- message == "Something happened"
- }
+ StepVerifier.create(act)
+ .expectErrorMatches { t ->
+ t instanceof JsonRpcException && t.error.code == RpcResponseError.CODE_UPSTREAM_INVALID_RESPONSE
+ }
+ .verify(Duration.ofSeconds(1))
}
}
diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/ResponseRpcParserSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/ResponseRpcParserSpec.groovy
index a7e23694..f2f3542d 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/ResponseRpcParserSpec.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/ResponseRpcParserSpec.groovy
@@ -223,27 +223,4 @@ class ResponseRpcParserSpec extends Specification {
!act.hasResult()
}
- def "Keep provided ID even if JSON is not full"() {
- setup:
- def json = '{"jsonrpc": "2.0", "id": 1}'
- when:
- def act = parser.parse(json.getBytes())
- then:
- act.id.asNumber() == 1
- act.error != null
- act.hasError()
- !act.hasResult()
- }
-
- def "Keep provided ID even if JSON is broken"() {
- setup:
- def json = '{"jsonrpc": "2.0", "id": 101, "resu'
- when:
- def act = parser.parse(json.getBytes())
- then:
- act.id.asNumber() == 101
- act.error != null
- act.hasError()
- !act.hasResult()
- }
}
diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParserSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParserSpec.groovy
index f2dab053..460e5e02 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParserSpec.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/rpcclient/ResponseWSParserSpec.groovy
@@ -103,24 +103,4 @@ class ResponseWSParserSpec extends Specification {
act.error == null
new String(act.value) == "null"
}
-
- def "Keep provided ID even if JSON is not full"() {
- setup:
- def json = '{"jsonrpc": "2.0", "id": 1}'
- when:
- def act = parser.parse(json.getBytes())
- then:
- act.id.asNumber() == 1
- act.error != null
- }
-
- def "Keep provided ID even if JSON is broken"() {
- setup:
- def json = '{"jsonrpc": "2.0", "id": 101, "resu'
- when:
- def act = parser.parse(json.getBytes())
- then:
- act.id.asNumber() == 101
- act.error != null
- }
}
diff --git a/src/test/groovy/io/emeraldpay/dshackle/upstream/signature/EcdsaSignerSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/upstream/signature/EcdsaSignerSpec.groovy
index 98a8acbb..668b5abc 100644
--- a/src/test/groovy/io/emeraldpay/dshackle/upstream/signature/EcdsaSignerSpec.groovy
+++ b/src/test/groovy/io/emeraldpay/dshackle/upstream/signature/EcdsaSignerSpec.groovy
@@ -49,7 +49,7 @@ class EcdsaSignerSpec extends Specification {
setup:
def conf = new SignatureConfig()
conf.enabled = true
- conf.privateKey = "src/test/resources/signer/test_key"
+ conf.privateKey = "testing/dshackle/test_key"
def signer = new ResponseSignerFactory(conf).getObject() as EcdsaSigner
// To verify the test, check the hash of test key above:
@@ -114,7 +114,7 @@ class EcdsaSignerSpec extends Specification {
def conf = new SignatureConfig()
conf.enabled = true
- conf.privateKey = "src/test/resources/signer/test_key"
+ conf.privateKey = "testing/dshackle/test_key"
def factory = new ResponseSignerFactory(conf)
def sk = factory.readKey(conf.algorithm, conf.privateKey).first
diff --git a/src/test/resources/configs/cache-redis-disabled.yaml b/src/test/resources/cache-redis-disabled.yaml
similarity index 100%
rename from src/test/resources/configs/cache-redis-disabled.yaml
rename to src/test/resources/cache-redis-disabled.yaml
diff --git a/src/test/resources/configs/cache-redis-full.yaml b/src/test/resources/cache-redis-full.yaml
similarity index 100%
rename from src/test/resources/configs/cache-redis-full.yaml
rename to src/test/resources/cache-redis-full.yaml
diff --git a/src/test/resources/configs/upstreams-priority.yaml b/src/test/resources/configs/upstreams-priority.yaml
deleted file mode 100644
index 2a165754..00000000
--- a/src/test/resources/configs/upstreams-priority.yaml
+++ /dev/null
@@ -1,26 +0,0 @@
-version: v1
-
-upstreams:
-
- - id: local
- chain: ethereum
- priority: 100
- connection:
- ethereum:
- rpc:
- url: "http://localhost:8545"
-
- - id: infura
- chain: ethereum
- role: fallback
- priority: 50
- connection:
- ethereum:
- rpc:
- url: "https://mainnet.infura.io/v3/fa28c968191849c1aff541ad1d8511f2"
-
- - id: remote
- priority: 75
- connection:
- grpc:
- host: "10.2.0.15"
\ No newline at end of file
diff --git a/src/test/resources/configs/dshackle-full.yaml b/src/test/resources/dshackle-full.yaml
similarity index 100%
rename from src/test/resources/configs/dshackle-full.yaml
rename to src/test/resources/dshackle-full.yaml
diff --git a/src/test/resources/configs/dshackle-health-1.yaml b/src/test/resources/dshackle-health-1.yaml
similarity index 100%
rename from src/test/resources/configs/dshackle-health-1.yaml
rename to src/test/resources/dshackle-health-1.yaml
diff --git a/src/test/resources/configs/dshackle-health-2.yaml b/src/test/resources/dshackle-health-2.yaml
similarity index 100%
rename from src/test/resources/configs/dshackle-health-2.yaml
rename to src/test/resources/dshackle-health-2.yaml
diff --git a/src/test/resources/configs/dshackle-health-empty.yaml b/src/test/resources/dshackle-health-empty.yaml
similarity index 100%
rename from src/test/resources/configs/dshackle-health-empty.yaml
rename to src/test/resources/dshackle-health-empty.yaml
diff --git a/src/test/resources/configs/dshackle-monitoring-basic.yaml b/src/test/resources/dshackle-monitoring-basic.yaml
similarity index 100%
rename from src/test/resources/configs/dshackle-monitoring-basic.yaml
rename to src/test/resources/dshackle-monitoring-basic.yaml
diff --git a/src/test/resources/configs/dshackle-proxy-basic.yaml b/src/test/resources/dshackle-proxy-basic.yaml
similarity index 100%
rename from src/test/resources/configs/dshackle-proxy-basic.yaml
rename to src/test/resources/dshackle-proxy-basic.yaml
diff --git a/src/test/resources/configs/dshackle-proxy-max.yaml b/src/test/resources/dshackle-proxy-max.yaml
similarity index 100%
rename from src/test/resources/configs/dshackle-proxy-max.yaml
rename to src/test/resources/dshackle-proxy-max.yaml
diff --git a/src/test/resources/configs/dshackle-proxy-no-ws.yaml b/src/test/resources/dshackle-proxy-no-ws.yaml
similarity index 100%
rename from src/test/resources/configs/dshackle-proxy-no-ws.yaml
rename to src/test/resources/dshackle-proxy-no-ws.yaml
diff --git a/src/test/resources/configs/dshackle-proxy-two.yaml b/src/test/resources/dshackle-proxy-two.yaml
similarity index 100%
rename from src/test/resources/configs/dshackle-proxy-two.yaml
rename to src/test/resources/dshackle-proxy-two.yaml
diff --git a/src/test/resources/configs/upstreams-basic.yaml b/src/test/resources/upstreams-basic.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-basic.yaml
rename to src/test/resources/upstreams-basic.yaml
diff --git a/src/test/resources/configs/upstreams-bitcoin-esplora.yaml b/src/test/resources/upstreams-bitcoin-esplora.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-bitcoin-esplora.yaml
rename to src/test/resources/upstreams-bitcoin-esplora.yaml
diff --git a/src/test/resources/configs/upstreams-bitcoin.yaml b/src/test/resources/upstreams-bitcoin.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-bitcoin.yaml
rename to src/test/resources/upstreams-bitcoin.yaml
diff --git a/src/test/resources/configs/upstreams-ds.yaml b/src/test/resources/upstreams-ds.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-ds.yaml
rename to src/test/resources/upstreams-ds.yaml
diff --git a/src/test/resources/configs/upstreams-extra.yaml b/src/test/resources/upstreams-extra.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-extra.yaml
rename to src/test/resources/upstreams-extra.yaml
diff --git a/src/test/resources/configs/upstreams-labels.yaml b/src/test/resources/upstreams-labels.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-labels.yaml
rename to src/test/resources/upstreams-labels.yaml
diff --git a/src/test/resources/configs/upstreams-methods-quorum.yaml b/src/test/resources/upstreams-methods-quorum.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-methods-quorum.yaml
rename to src/test/resources/upstreams-methods-quorum.yaml
diff --git a/src/test/resources/configs/upstreams-methods.yaml b/src/test/resources/upstreams-methods.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-methods.yaml
rename to src/test/resources/upstreams-methods.yaml
diff --git a/src/test/resources/configs/upstreams-no-defaults.yaml b/src/test/resources/upstreams-no-defaults.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-no-defaults.yaml
rename to src/test/resources/upstreams-no-defaults.yaml
diff --git a/src/test/resources/configs/upstreams-no-id.yaml b/src/test/resources/upstreams-no-id.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-no-id.yaml
rename to src/test/resources/upstreams-no-id.yaml
diff --git a/src/test/resources/configs/upstreams-options.yaml b/src/test/resources/upstreams-options.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-options.yaml
rename to src/test/resources/upstreams-options.yaml
diff --git a/src/test/resources/configs/upstreams-roles-2.yaml b/src/test/resources/upstreams-roles-2.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-roles-2.yaml
rename to src/test/resources/upstreams-roles-2.yaml
diff --git a/src/test/resources/configs/upstreams-roles-invalid.yaml b/src/test/resources/upstreams-roles-invalid.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-roles-invalid.yaml
rename to src/test/resources/upstreams-roles-invalid.yaml
diff --git a/src/test/resources/configs/upstreams-roles.yaml b/src/test/resources/upstreams-roles.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-roles.yaml
rename to src/test/resources/upstreams-roles.yaml
diff --git a/src/test/resources/configs/upstreams-ws-full.yaml b/src/test/resources/upstreams-ws-full.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-ws-full.yaml
rename to src/test/resources/upstreams-ws-full.yaml
diff --git a/src/test/resources/configs/upstreams-ws-only.yaml b/src/test/resources/upstreams-ws-only.yaml
similarity index 100%
rename from src/test/resources/configs/upstreams-ws-only.yaml
rename to src/test/resources/upstreams-ws-only.yaml
diff --git a/testing/.gitignore b/testing/.gitignore
new file mode 100644
index 00000000..e5a210d4
--- /dev/null
+++ b/testing/.gitignore
@@ -0,0 +1,3 @@
+gradlew
+gradlew.bat
+gradle/
\ No newline at end of file
diff --git a/testing/README.adoc b/testing/README.adoc
new file mode 100644
index 00000000..8027d983
--- /dev/null
+++ b/testing/README.adoc
@@ -0,0 +1,104 @@
+= Dshackle Integration Testing
+
+Test for a correct work using a simulation of different upstream behaviours and checking the Dshackle against it.
+
+.Modules:
+- `dshackle` - Dshackle config used for testing environment
+- `simple-upstream` - upstream simulator with predefines responses
+- `trial` - actual tests
+
+== How to run
+
+NOTE: Run commands in project's root dir
+
+NOTE: You need to open multiple consoles, since each command runs on foreground
+
+== Basic Test
+
+Makes some basic checks with a mocked upstreams
+
+=== Run upstream emulator
+
+.Runs two instances of RPC server
+[source,bash]
+----
+DSHACKLE_TESTUP_PORT=18545 ./gradlew -p testing/simple-upstream run
+DSHACKLE_TESTUP_PORT=18546 ./gradlew -p testing/simple-upstream run
+----
+
+=== Run Dshackle
+
+[source,bash]
+----
+./gradlew run --args="--configPath=./testing/dshackle/dshackle-basic.yaml"
+----
+
+=== Run tests
+
+[source,bash]
+----
+./testing/trial/gradlew -p testing/trial -PdshackleTrialMode=basic cleanTest test
+----
+
+== Real (Mainnet) Test
+
+=== Run Dshackle
+
+First you have to provide configuration to an upstream connected to the mainnet.
+It may be a local Geth instance, or an other provider, like Infura (note that the test config doesn't support TLS or other auth)
+
+[source,bash]
+----
+export DSHACKLE_TEST_ETH1_RPC=
+export DSHACKLE_TEST_ETH1_WS=
+export DSHACKLE_TEST_ETH1_WSORIGIN=
+----
+
+For a local Geth it would be:
+
+[source,bash]
+----
+export DSHACKLE_TEST_ETH1_RPC=http://127.0.0.1:8545
+export DSHACKLE_TEST_ETH1_WS=ws://127.0.0.1:8546
+export DSHACKLE_TEST_ETH1_WSORIGIN=http://127.0.0.1:8546
+----
+
+Finally, run the Dshackle instance:
+
+[source,bash]
+----
+./gradlew run --args="--configPath=./testing/dshackle/dshackle-real.yaml"
+----
+
+=== Run tests
+
+[source,bash]
+----
+./testing/trial/gradlew -p testing/trial -PdshackleTrialMode=real cleanTest test
+----
+
+== Real (Rinkeby) Test
+
+Runs requests against a Rinkeby Testnet.
+Configuration is similar to the Mainnet config, but instead of `ETH1` is has `RINKEBY`.
+
+[source,bash]
+----
+export DSHACKLE_TEST_RINKEBY_RPC=
+export DSHACKLE_TEST_RINKEBY_WS=
+export DSHACKLE_TEST_RINKEBY_WSORIGIN=
+----
+
+And run the Dshackle instance:
+
+[source,bash]
+----
+./gradlew run --args="--configPath=./testing/dshackle/dshackle-rinkeby.yaml"
+----
+
+=== Run tests
+
+[source,bash]
+----
+./testing/trial/gradlew -p testing/trial -PdshackleTrialMode=rinkeby cleanTest test
+----
diff --git a/testing/dshackle/README.adoc b/testing/dshackle/README.adoc
new file mode 100644
index 00000000..46dc3bc3
--- /dev/null
+++ b/testing/dshackle/README.adoc
@@ -0,0 +1,14 @@
+= Testing Server
+
+.Run from project root directory:
+[source,bash]
+----
+./gradlew run --args="--configPath=./testing/dshackle/dshackle-mock.yaml"
+----
+
+or
+
+[source,bash]
+----
+./gradlew run --args="--configPath=./testing/dshackle/dshackle-real.yaml"
+----
diff --git a/testing/dshackle/dshackle-basic.yaml b/testing/dshackle/dshackle-basic.yaml
new file mode 100644
index 00000000..56515c78
--- /dev/null
+++ b/testing/dshackle/dshackle-basic.yaml
@@ -0,0 +1,50 @@
+version: v1
+port: 12448
+
+tls:
+ enabled: false
+
+cluster:
+ upstreams:
+ - id: test-1
+ node-id: 1
+ chain: ethereum
+ methods:
+ enabled:
+ - name: debug_traceTransaction
+ - name: test_foo
+ options:
+ disable-validation: true
+ connection:
+ ethereum-pos:
+ execution:
+ rpc:
+ url: "http://localhost:18545"
+
+ - id: test-2
+ node-id: 2
+ chain: ethereum
+ options:
+ disable-validation: true
+ connection:
+ execution:
+ ethereum-pos:
+ rpc:
+ url: "http://localhost:18546"
+
+cache:
+ redis:
+ enabled: false
+
+signed-response:
+ enabled: true
+ algorithm: SECP256K1
+ private-key: "./test_key"
+
+proxy:
+ port: 18080
+ tls:
+ enabled: false
+ routes:
+ - id: eth
+ blockchain: ethereum
\ No newline at end of file
diff --git a/testing/dshackle/dshackle-real.yaml b/testing/dshackle/dshackle-real.yaml
new file mode 100644
index 00000000..a8e289f0
--- /dev/null
+++ b/testing/dshackle/dshackle-real.yaml
@@ -0,0 +1,30 @@
+version: v1
+port: 12448
+
+tls:
+ enabled: false
+
+cluster:
+ upstreams:
+ - id: eth-1
+ chain: ethereum
+ connection:
+ ethereum:
+ rpc:
+ url: "${DSHACKLE_TEST_ETH1_RPC}"
+ ws:
+ url: "${DSHACKLE_TEST_ETH1_WS}"
+ origin: "${DSHACKLE_TEST_ETH1_WSORIGIN}"
+
+cache:
+ redis:
+ enabled: false
+
+proxy:
+ port: 18081
+ preserve-batch-order: true
+ tls:
+ enabled: false
+ routes:
+ - id: eth
+ blockchain: ethereum
\ No newline at end of file
diff --git a/testing/dshackle/dshackle-rinkeby.yaml b/testing/dshackle/dshackle-rinkeby.yaml
new file mode 100644
index 00000000..e0e1d3a8
--- /dev/null
+++ b/testing/dshackle/dshackle-rinkeby.yaml
@@ -0,0 +1,26 @@
+version: v1
+port: 12448
+
+tls:
+ enabled: false
+
+cluster:
+ upstreams:
+ - id: rinkeby-1
+ chain: rinkeby
+ connection:
+ ethereum:
+ rpc:
+ url: "${DSHACKLE_TEST_RINKEBY_RPC}"
+
+cache:
+ redis:
+ enabled: false
+
+proxy:
+ port: 18081
+ tls:
+ enabled: false
+ routes:
+ - id: rinkeby
+ blockchain: rinkeby
\ No newline at end of file
diff --git a/src/test/resources/signer/test_key b/testing/dshackle/test_key
similarity index 100%
rename from src/test/resources/signer/test_key
rename to testing/dshackle/test_key
diff --git a/src/test/resources/signer/test_key.pub b/testing/dshackle/test_key.pub
similarity index 100%
rename from src/test/resources/signer/test_key.pub
rename to testing/dshackle/test_key.pub
diff --git a/testing/simple-upstream/README.adoc b/testing/simple-upstream/README.adoc
new file mode 100644
index 00000000..caab006c
--- /dev/null
+++ b/testing/simple-upstream/README.adoc
@@ -0,0 +1,6 @@
+= Simple Testing Upstream
+
+.Run from current directory
+----
+./gradlew run
+----
\ No newline at end of file
diff --git a/testing/simple-upstream/build.gradle b/testing/simple-upstream/build.gradle
new file mode 100644
index 00000000..72b72e50
--- /dev/null
+++ b/testing/simple-upstream/build.gradle
@@ -0,0 +1,22 @@
+plugins {
+ id 'java'
+ id 'groovy'
+ id 'idea'
+ id 'application'
+}
+
+repositories {
+ mavenLocal()
+ mavenCentral()
+}
+
+dependencies {
+ implementation "com.sparkjava:spark-core:2.9.1"
+ implementation "org.codehaus.groovy:groovy:3.0.4"
+ implementation "com.fasterxml.jackson.core:jackson-core:2.9.8"
+ implementation "com.fasterxml.jackson.core:jackson-databind:2.9.8"
+}
+
+application {
+ mainClassName = 'testing.SimpleUpstream'
+}
\ No newline at end of file
diff --git a/testing/simple-upstream/settings.gradle b/testing/simple-upstream/settings.gradle
new file mode 100644
index 00000000..faa039c6
--- /dev/null
+++ b/testing/simple-upstream/settings.gradle
@@ -0,0 +1 @@
+rootProject.name = 'dshackle-testing-simple-upstream'
\ No newline at end of file
diff --git a/testing/simple-upstream/src/main/groovy/testing/BlocksHandler.groovy b/testing/simple-upstream/src/main/groovy/testing/BlocksHandler.groovy
new file mode 100644
index 00000000..11526f97
--- /dev/null
+++ b/testing/simple-upstream/src/main/groovy/testing/BlocksHandler.groovy
@@ -0,0 +1,32 @@
+package testing
+
+import com.fasterxml.jackson.databind.ObjectMapper
+
+class BlocksHandler implements CallHandler {
+
+ ObjectMapper objectMapper
+ ResourceResponse resourceResponse
+
+ BlocksHandler(ObjectMapper objectMapper) {
+ this.objectMapper = objectMapper
+ this.resourceResponse = new ResourceResponse(objectMapper)
+ }
+
+ @Override
+ Result handle(String method, List