solution: refactor AccessLog handler
This commit is contained in:
@@ -51,20 +51,31 @@ class AccessHandler(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private fun <ReqT : Any, RespT : Any, E> process(
|
||||||
|
call: ServerCall<ReqT, RespT>,
|
||||||
|
headers: Metadata,
|
||||||
|
next: ServerCallHandler<ReqT, RespT>,
|
||||||
|
builder: EventsBuilder.RequestReply<E, ReqT, RespT>
|
||||||
|
): ServerCall.Listener<ReqT> {
|
||||||
|
builder.start(headers, call.attributes)
|
||||||
|
val callWrapper: ServerCall<ReqT, RespT> = StdCallResponse(
|
||||||
|
call, builder, accessLogWriter
|
||||||
|
)
|
||||||
|
return StdCallListener(
|
||||||
|
next.startCall(callWrapper, headers),
|
||||||
|
builder
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
@Suppress("UNCHECKED_CAST")
|
@Suppress("UNCHECKED_CAST")
|
||||||
private fun <ReqT : Any, RespT : Any> processSubscribeHead(
|
private fun <ReqT : Any, RespT : Any> processSubscribeHead(
|
||||||
call: ServerCall<ReqT, RespT>,
|
call: ServerCall<ReqT, RespT>,
|
||||||
headers: Metadata,
|
headers: Metadata,
|
||||||
next: ServerCallHandler<ReqT, RespT>
|
next: ServerCallHandler<ReqT, RespT>
|
||||||
): ServerCall.Listener<ReqT> {
|
): ServerCall.Listener<ReqT> {
|
||||||
val builder = EventsBuilder.SubscribeHead()
|
return process(call, headers, next,
|
||||||
.start(headers, call.attributes)
|
EventsBuilder.SubscribeHead() as EventsBuilder.RequestReply<*, ReqT, RespT>
|
||||||
val callWrapper: ServerCall<ReqT, RespT> = OnSubscribeHeadResponse(
|
)
|
||||||
call as ServerCall<Common.Chain, BlockchainOuterClass.ChainHead>, builder, accessLogWriter) as ServerCall<ReqT, RespT>
|
|
||||||
return OnSubscribeHead(
|
|
||||||
next.startCall(callWrapper, headers) as ServerCall.Listener<Common.Chain>,
|
|
||||||
builder
|
|
||||||
) as ServerCall.Listener<ReqT>
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Suppress("UNCHECKED_CAST")
|
@Suppress("UNCHECKED_CAST")
|
||||||
@@ -74,14 +85,9 @@ class AccessHandler(
|
|||||||
next: ServerCallHandler<ReqT, RespT>,
|
next: ServerCallHandler<ReqT, RespT>,
|
||||||
subscribe: Boolean
|
subscribe: Boolean
|
||||||
): ServerCall.Listener<ReqT> {
|
): ServerCall.Listener<ReqT> {
|
||||||
val builder = EventsBuilder.SubscribeBalance(subscribe)
|
return process(call, headers, next,
|
||||||
.start(headers, call.attributes)
|
EventsBuilder.SubscribeBalance(subscribe) as EventsBuilder.RequestReply<*, ReqT, RespT>
|
||||||
val callWrapper: ServerCall<ReqT, RespT> = OnSubscribeBalanceResponse(
|
)
|
||||||
call as ServerCall<BlockchainOuterClass.BalanceRequest, BlockchainOuterClass.AddressBalance>, builder, accessLogWriter) as ServerCall<ReqT, RespT>
|
|
||||||
return OnSubscribeBalance(
|
|
||||||
next.startCall(callWrapper, headers) as ServerCall.Listener<BlockchainOuterClass.BalanceRequest>,
|
|
||||||
builder
|
|
||||||
) as ServerCall.Listener<ReqT>
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Suppress("UNCHECKED_CAST")
|
@Suppress("UNCHECKED_CAST")
|
||||||
@@ -90,14 +96,9 @@ class AccessHandler(
|
|||||||
headers: Metadata,
|
headers: Metadata,
|
||||||
next: ServerCallHandler<ReqT, RespT>
|
next: ServerCallHandler<ReqT, RespT>
|
||||||
): ServerCall.Listener<ReqT> {
|
): ServerCall.Listener<ReqT> {
|
||||||
val builder = EventsBuilder.TxStatus()
|
return process(call, headers, next,
|
||||||
.start(headers, call.attributes)
|
EventsBuilder.TxStatus() as EventsBuilder.RequestReply<*, ReqT, RespT>
|
||||||
val callWrapper: ServerCall<ReqT, RespT> = OnTxStatusResponse(
|
)
|
||||||
call as ServerCall<BlockchainOuterClass.TxStatusRequest, BlockchainOuterClass.TxStatus>, builder, accessLogWriter) as ServerCall<ReqT, RespT>
|
|
||||||
return OnSubscribeTxStatus(
|
|
||||||
next.startCall(callWrapper, headers) as ServerCall.Listener<BlockchainOuterClass.TxStatusRequest>,
|
|
||||||
builder
|
|
||||||
) as ServerCall.Listener<ReqT>
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Suppress("UNCHECKED_CAST")
|
@Suppress("UNCHECKED_CAST")
|
||||||
@@ -106,18 +107,9 @@ class AccessHandler(
|
|||||||
headers: Metadata,
|
headers: Metadata,
|
||||||
next: ServerCallHandler<ReqT, RespT>
|
next: ServerCallHandler<ReqT, RespT>
|
||||||
): ServerCall.Listener<ReqT> {
|
): ServerCall.Listener<ReqT> {
|
||||||
val builder = EventsBuilder.NativeCall()
|
return process(call, headers, next,
|
||||||
.start(headers, call.attributes)
|
EventsBuilder.NativeCall() as EventsBuilder.RequestReply<*, ReqT, RespT>
|
||||||
|
)
|
||||||
val callWrapper: ServerCall<ReqT, RespT> = OnNativeCallResponse(
|
|
||||||
call as ServerCall<BlockchainOuterClass.NativeCallRequest, BlockchainOuterClass.NativeCallReplyItem>,
|
|
||||||
builder,
|
|
||||||
accessLogWriter
|
|
||||||
) as ServerCall<ReqT, RespT>
|
|
||||||
return OnNativeCall(
|
|
||||||
next.startCall(callWrapper, headers) as ServerCall.Listener<BlockchainOuterClass.NativeCallRequest>,
|
|
||||||
builder
|
|
||||||
) as ServerCall.Listener<ReqT>
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Suppress("UNCHECKED_CAST")
|
@Suppress("UNCHECKED_CAST")
|
||||||
@@ -126,17 +118,9 @@ class AccessHandler(
|
|||||||
headers: Metadata,
|
headers: Metadata,
|
||||||
next: ServerCallHandler<ReqT, RespT>
|
next: ServerCallHandler<ReqT, RespT>
|
||||||
): ServerCall.Listener<ReqT> {
|
): ServerCall.Listener<ReqT> {
|
||||||
val builder = EventsBuilder.Describe()
|
return process(call, headers, next,
|
||||||
.start(headers, call.attributes)
|
EventsBuilder.Describe() as EventsBuilder.RequestReply<*, ReqT, RespT>
|
||||||
|
)
|
||||||
val callWrapper: ServerCall<ReqT, RespT> = OnDescribeResponse(
|
|
||||||
call as ServerCall<BlockchainOuterClass.DescribeRequest, BlockchainOuterClass.DescribeResponse>,
|
|
||||||
builder,
|
|
||||||
accessLogWriter
|
|
||||||
) as ServerCall<ReqT, RespT>
|
|
||||||
return OnDescribeRequest(
|
|
||||||
next.startCall(callWrapper, headers) as ServerCall.Listener<BlockchainOuterClass.DescribeRequest>,
|
|
||||||
builder) as ServerCall.Listener<ReqT>
|
|
||||||
}
|
}
|
||||||
|
|
||||||
@Suppress("UNCHECKED_CAST")
|
@Suppress("UNCHECKED_CAST")
|
||||||
@@ -145,115 +129,32 @@ class AccessHandler(
|
|||||||
headers: Metadata,
|
headers: Metadata,
|
||||||
next: ServerCallHandler<ReqT, RespT>
|
next: ServerCallHandler<ReqT, RespT>
|
||||||
): ServerCall.Listener<ReqT> {
|
): ServerCall.Listener<ReqT> {
|
||||||
val builder = EventsBuilder.Status()
|
return process(call, headers, next,
|
||||||
.start(headers, call.attributes)
|
EventsBuilder.Status() as EventsBuilder.RequestReply<*, ReqT, RespT>
|
||||||
|
)
|
||||||
val callWrapper: ServerCall<ReqT, RespT> = OnStatusResponse(
|
|
||||||
call as ServerCall<BlockchainOuterClass.StatusRequest, BlockchainOuterClass.ChainStatus>,
|
|
||||||
builder,
|
|
||||||
accessLogWriter
|
|
||||||
) as ServerCall<ReqT, RespT>
|
|
||||||
return OnStatusRequest(
|
|
||||||
next.startCall(callWrapper, headers) as ServerCall.Listener<BlockchainOuterClass.StatusRequest>,
|
|
||||||
builder) as ServerCall.Listener<ReqT>
|
|
||||||
}
|
}
|
||||||
|
|
||||||
class OnSubscribeHead(
|
open class StdCallListener<Req, EB : EventsBuilder.RequestReply<*, Req, *>>(
|
||||||
val next: ServerCall.Listener<Common.Chain>,
|
val next: ServerCall.Listener<Req>,
|
||||||
val builder: EventsBuilder.SubscribeHead
|
val builder: EB
|
||||||
) : ForwardingServerCallListener<Common.Chain>() {
|
) : ForwardingServerCallListener<Req>() {
|
||||||
|
|
||||||
override fun onMessage(message: Common.Chain) {
|
override fun onMessage(message: Req) {
|
||||||
val chainId = message.type.number
|
builder.onRequest(message)
|
||||||
builder.withChain(chainId)
|
|
||||||
super.onMessage(message)
|
super.onMessage(message)
|
||||||
}
|
}
|
||||||
|
|
||||||
override fun delegate(): ServerCall.Listener<Common.Chain> {
|
override fun delegate(): ServerCall.Listener<Req> {
|
||||||
return next
|
return next
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
class OnSubscribeBalance(
|
open class StdCallResponse<ReqT : Any, RespT : Any, EB : EventsBuilder.RequestReply<*, ReqT, RespT>>(
|
||||||
val next: ServerCall.Listener<BlockchainOuterClass.BalanceRequest>,
|
val next: ServerCall<ReqT, RespT>,
|
||||||
val builder: EventsBuilder.SubscribeBalance
|
val builder: EB,
|
||||||
) : ForwardingServerCallListener<BlockchainOuterClass.BalanceRequest>() {
|
val accessLogWriter: AccessLogWriter
|
||||||
|
|
||||||
override fun onMessage(message: BlockchainOuterClass.BalanceRequest) {
|
|
||||||
builder.withRequest(message)
|
|
||||||
super.onMessage(message)
|
|
||||||
}
|
|
||||||
|
|
||||||
override fun delegate(): ServerCall.Listener<BlockchainOuterClass.BalanceRequest> {
|
|
||||||
return next
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
class OnSubscribeTxStatus(
|
|
||||||
val next: ServerCall.Listener<BlockchainOuterClass.TxStatusRequest>,
|
|
||||||
val builder: EventsBuilder.TxStatus
|
|
||||||
) : ForwardingServerCallListener<BlockchainOuterClass.TxStatusRequest>() {
|
|
||||||
|
|
||||||
override fun onMessage(message: BlockchainOuterClass.TxStatusRequest) {
|
|
||||||
builder.withRequest(message)
|
|
||||||
super.onMessage(message)
|
|
||||||
}
|
|
||||||
|
|
||||||
override fun delegate(): ServerCall.Listener<BlockchainOuterClass.TxStatusRequest> {
|
|
||||||
return next
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
class OnNativeCall(
|
|
||||||
val next: ServerCall.Listener<BlockchainOuterClass.NativeCallRequest>,
|
|
||||||
val builder: EventsBuilder.NativeCall
|
|
||||||
) : ForwardingServerCallListener<BlockchainOuterClass.NativeCallRequest>() {
|
|
||||||
|
|
||||||
override fun onMessage(message: BlockchainOuterClass.NativeCallRequest) {
|
|
||||||
val chain = message.chain
|
|
||||||
builder.withChain(chain.number)
|
|
||||||
message.itemsList.forEach { item ->
|
|
||||||
builder.onRequest(item)
|
|
||||||
}
|
|
||||||
super.onMessage(message)
|
|
||||||
}
|
|
||||||
|
|
||||||
override fun delegate(): ServerCall.Listener<BlockchainOuterClass.NativeCallRequest> {
|
|
||||||
return next
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
class OnDescribeRequest(
|
|
||||||
val next: ServerCall.Listener<BlockchainOuterClass.DescribeRequest>,
|
|
||||||
val builder: EventsBuilder.Describe
|
|
||||||
) : ForwardingServerCallListener<BlockchainOuterClass.DescribeRequest>() {
|
|
||||||
|
|
||||||
override fun onMessage(message: BlockchainOuterClass.DescribeRequest) {
|
|
||||||
super.onMessage(message)
|
|
||||||
}
|
|
||||||
|
|
||||||
override fun delegate(): ServerCall.Listener<BlockchainOuterClass.DescribeRequest> {
|
|
||||||
return next
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
class OnStatusRequest(
|
|
||||||
val next: ServerCall.Listener<BlockchainOuterClass.StatusRequest>,
|
|
||||||
val builder: EventsBuilder.Status
|
|
||||||
) : ForwardingServerCallListener<BlockchainOuterClass.StatusRequest>() {
|
|
||||||
|
|
||||||
override fun onMessage(message: BlockchainOuterClass.StatusRequest) {
|
|
||||||
super.onMessage(message)
|
|
||||||
}
|
|
||||||
|
|
||||||
override fun delegate(): ServerCall.Listener<BlockchainOuterClass.StatusRequest> {
|
|
||||||
return next
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
abstract class BaseCallResponse<ReqT : Any, RespT : Any>(
|
|
||||||
val next: ServerCall<ReqT, RespT>
|
|
||||||
) : ForwardingServerCall<ReqT, RespT>() {
|
) : ForwardingServerCall<ReqT, RespT>() {
|
||||||
|
|
||||||
override fun getMethodDescriptor(): MethodDescriptor<ReqT, RespT> {
|
override fun getMethodDescriptor(): MethodDescriptor<ReqT, RespT> {
|
||||||
return next.methodDescriptor
|
return next.methodDescriptor
|
||||||
}
|
}
|
||||||
@@ -264,85 +165,10 @@ class AccessHandler(
|
|||||||
|
|
||||||
override fun sendMessage(message: RespT) {
|
override fun sendMessage(message: RespT) {
|
||||||
super.sendMessage(message)
|
super.sendMessage(message)
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
class OnNativeCallResponse(
|
|
||||||
next: ServerCall<BlockchainOuterClass.NativeCallRequest, BlockchainOuterClass.NativeCallReplyItem>,
|
|
||||||
val builder: EventsBuilder.NativeCall,
|
|
||||||
val accessLogWriter: AccessLogWriter
|
|
||||||
) : BaseCallResponse<BlockchainOuterClass.NativeCallRequest, BlockchainOuterClass.NativeCallReplyItem>(next) {
|
|
||||||
|
|
||||||
override fun sendMessage(message: BlockchainOuterClass.NativeCallReplyItem) {
|
|
||||||
accessLogWriter.submit(
|
accessLogWriter.submit(
|
||||||
builder.onReply(message)
|
builder.onReply(message)!!
|
||||||
)
|
)
|
||||||
super.sendMessage(message)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
class OnSubscribeHeadResponse(
|
|
||||||
next: ServerCall<Common.Chain, BlockchainOuterClass.ChainHead>,
|
|
||||||
val builder: EventsBuilder.SubscribeHead,
|
|
||||||
val accessLogWriter: AccessLogWriter
|
|
||||||
) : BaseCallResponse<Common.Chain, BlockchainOuterClass.ChainHead>(next) {
|
|
||||||
|
|
||||||
override fun sendMessage(message: BlockchainOuterClass.ChainHead) {
|
|
||||||
val event = builder.onReply(message)
|
|
||||||
accessLogWriter.submit(event)
|
|
||||||
super.sendMessage(message)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
class OnSubscribeBalanceResponse(
|
|
||||||
next: ServerCall<BlockchainOuterClass.BalanceRequest, BlockchainOuterClass.AddressBalance>,
|
|
||||||
val builder: EventsBuilder.SubscribeBalance,
|
|
||||||
val accessLogWriter: AccessLogWriter
|
|
||||||
) : BaseCallResponse<BlockchainOuterClass.BalanceRequest, BlockchainOuterClass.AddressBalance>(next) {
|
|
||||||
|
|
||||||
override fun sendMessage(message: BlockchainOuterClass.AddressBalance) {
|
|
||||||
val event = builder.onReply(message)
|
|
||||||
accessLogWriter.submit(event)
|
|
||||||
super.sendMessage(message)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
class OnTxStatusResponse(
|
|
||||||
next: ServerCall<BlockchainOuterClass.TxStatusRequest, BlockchainOuterClass.TxStatus>,
|
|
||||||
val builder: EventsBuilder.TxStatus,
|
|
||||||
val accessLogWriter: AccessLogWriter
|
|
||||||
) : BaseCallResponse<BlockchainOuterClass.TxStatusRequest, BlockchainOuterClass.TxStatus>(next) {
|
|
||||||
|
|
||||||
override fun sendMessage(message: BlockchainOuterClass.TxStatus) {
|
|
||||||
val event = builder.onReply(message)
|
|
||||||
accessLogWriter.submit(event)
|
|
||||||
super.sendMessage(message)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
class OnDescribeResponse(
|
|
||||||
next: ServerCall<BlockchainOuterClass.DescribeRequest, BlockchainOuterClass.DescribeResponse>,
|
|
||||||
val builder: EventsBuilder.Describe,
|
|
||||||
val accessLogWriter: AccessLogWriter
|
|
||||||
) : BaseCallResponse<BlockchainOuterClass.DescribeRequest, BlockchainOuterClass.DescribeResponse>(next) {
|
|
||||||
|
|
||||||
override fun sendMessage(message: BlockchainOuterClass.DescribeResponse) {
|
|
||||||
val event = builder.onReply()
|
|
||||||
accessLogWriter.submit(event)
|
|
||||||
super.sendMessage(message)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
class OnStatusResponse(
|
|
||||||
next: ServerCall<BlockchainOuterClass.StatusRequest, BlockchainOuterClass.ChainStatus>,
|
|
||||||
val builder: EventsBuilder.Status,
|
|
||||||
val accessLogWriter: AccessLogWriter
|
|
||||||
) : BaseCallResponse<BlockchainOuterClass.StatusRequest, BlockchainOuterClass.ChainStatus>(next) {
|
|
||||||
|
|
||||||
override fun sendMessage(message: BlockchainOuterClass.ChainStatus) {
|
|
||||||
val event = builder.onReply(message)
|
|
||||||
accessLogWriter.submit(event)
|
|
||||||
super.sendMessage(message)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
@@ -16,6 +16,7 @@
|
|||||||
package io.emeraldpay.dshackle.monitoring.accesslog
|
package io.emeraldpay.dshackle.monitoring.accesslog
|
||||||
|
|
||||||
import io.emeraldpay.api.proto.BlockchainOuterClass
|
import io.emeraldpay.api.proto.BlockchainOuterClass
|
||||||
|
import io.emeraldpay.api.proto.Common
|
||||||
import io.emeraldpay.grpc.Chain
|
import io.emeraldpay.grpc.Chain
|
||||||
import io.grpc.Attributes
|
import io.grpc.Attributes
|
||||||
import io.grpc.Grpc
|
import io.grpc.Grpc
|
||||||
@@ -33,7 +34,16 @@ class EventsBuilder {
|
|||||||
private val log = LoggerFactory.getLogger(EventsBuilder::class.java)
|
private val log = LoggerFactory.getLogger(EventsBuilder::class.java)
|
||||||
}
|
}
|
||||||
|
|
||||||
abstract class Base<T>() {
|
interface StartingRequest {
|
||||||
|
fun start(metadata: Metadata, attributes: Attributes)
|
||||||
|
}
|
||||||
|
|
||||||
|
interface RequestReply<E, Req, Resp> : StartingRequest {
|
||||||
|
fun onRequest(msg: Req)
|
||||||
|
fun onReply(msg: Resp): E
|
||||||
|
}
|
||||||
|
|
||||||
|
abstract class Base<T>() : StartingRequest {
|
||||||
companion object {
|
companion object {
|
||||||
private val remoteIpKeys = listOf(
|
private val remoteIpKeys = listOf(
|
||||||
Metadata.Key.of("x-real-ip", Metadata.ASCII_STRING_MARSHALLER),
|
Metadata.Key.of("x-real-ip", Metadata.ASCII_STRING_MARSHALLER),
|
||||||
@@ -82,9 +92,9 @@ class EventsBuilder {
|
|||||||
.trim()
|
.trim()
|
||||||
}
|
}
|
||||||
|
|
||||||
abstract protected fun getT(): T
|
protected abstract fun getT(): T
|
||||||
|
|
||||||
fun start(metadata: Metadata, attributes: Attributes): T {
|
override fun start(metadata: Metadata, attributes: Attributes) {
|
||||||
val userAgent = metadata.get(Metadata.Key.of("user-agent", Metadata.ASCII_STRING_MARSHALLER))
|
val userAgent = metadata.get(Metadata.Key.of("user-agent", Metadata.ASCII_STRING_MARSHALLER))
|
||||||
?.let(this@Base::clean)
|
?.let(this@Base::clean)
|
||||||
?: ""
|
?: ""
|
||||||
@@ -108,7 +118,6 @@ class EventsBuilder {
|
|||||||
ip = ip,
|
ip = ip,
|
||||||
userAgent = userAgent
|
userAgent = userAgent
|
||||||
))
|
))
|
||||||
return getT()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fun withChain(chain: Int): T {
|
fun withChain(chain: Int): T {
|
||||||
@@ -118,21 +127,31 @@ class EventsBuilder {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
class SubscribeHead() : Base<SubscribeHead>() {
|
class SubscribeHead() :
|
||||||
|
Base<SubscribeHead>(),
|
||||||
|
RequestReply<Events.SubscribeHead, Common.Chain, BlockchainOuterClass.ChainHead> {
|
||||||
|
|
||||||
private var index = 0
|
private var index = 0
|
||||||
|
|
||||||
override fun getT(): SubscribeHead {
|
override fun getT(): SubscribeHead {
|
||||||
return this
|
return this
|
||||||
}
|
}
|
||||||
|
|
||||||
fun onReply(resp: BlockchainOuterClass.ChainHead): Events.SubscribeHead {
|
override fun onRequest(msg: Common.Chain) {
|
||||||
|
withChain(msg.type.number)
|
||||||
|
}
|
||||||
|
|
||||||
|
override fun onReply(msg: BlockchainOuterClass.ChainHead): Events.SubscribeHead {
|
||||||
return Events.SubscribeHead(
|
return Events.SubscribeHead(
|
||||||
chain, UUID.randomUUID(), requestDetails, index++
|
chain, UUID.randomUUID(), requestDetails, index++
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
class SubscribeBalance(val subscribe: Boolean) : Base<SubscribeBalance>() {
|
class SubscribeBalance(val subscribe: Boolean) :
|
||||||
|
Base<SubscribeBalance>(),
|
||||||
|
RequestReply<Events.SubscribeBalance, BlockchainOuterClass.BalanceRequest, BlockchainOuterClass.AddressBalance> {
|
||||||
|
|
||||||
private var index = 0
|
private var index = 0
|
||||||
private var balanceRequest: Events.BalanceRequest? = null
|
private var balanceRequest: Events.BalanceRequest? = null
|
||||||
|
|
||||||
@@ -140,39 +159,40 @@ class EventsBuilder {
|
|||||||
return this
|
return this
|
||||||
}
|
}
|
||||||
|
|
||||||
fun withRequest(req: BlockchainOuterClass.BalanceRequest): SubscribeBalance {
|
override fun onRequest(msg: BlockchainOuterClass.BalanceRequest) {
|
||||||
balanceRequest = Events.BalanceRequest(
|
balanceRequest = Events.BalanceRequest(
|
||||||
req.asset.code.toUpperCase(),
|
msg.asset.code.toUpperCase(),
|
||||||
req.address.addrTypeCase.name
|
msg.address.addrTypeCase.name
|
||||||
)
|
)
|
||||||
return this
|
|
||||||
}
|
}
|
||||||
|
|
||||||
fun onReply(resp: BlockchainOuterClass.AddressBalance): Events.SubscribeBalance {
|
override fun onReply(msg: BlockchainOuterClass.AddressBalance): Events.SubscribeBalance {
|
||||||
if (balanceRequest == null) {
|
if (balanceRequest == null) {
|
||||||
throw IllegalStateException("Request is not initialized")
|
throw IllegalStateException("Request is not initialized")
|
||||||
}
|
}
|
||||||
val addressBalance = Events.AddressBalance(resp.asset.code, resp.address.address)
|
val addressBalance = Events.AddressBalance(msg.asset.code, msg.address.address)
|
||||||
val chain = Chain.byId(resp.asset.chain.number)
|
val chain = Chain.byId(msg.asset.chain.number)
|
||||||
return Events.SubscribeBalance(
|
return Events.SubscribeBalance(
|
||||||
chain, UUID.randomUUID(), subscribe, requestDetails, balanceRequest!!, addressBalance, index++
|
chain, UUID.randomUUID(), subscribe, requestDetails, balanceRequest!!, addressBalance, index++
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
class TxStatus() : Base<TxStatus>() {
|
class TxStatus() :
|
||||||
|
Base<TxStatus>(),
|
||||||
|
RequestReply<Events.TxStatus, BlockchainOuterClass.TxStatusRequest, BlockchainOuterClass.TxStatus> {
|
||||||
private var index = 0
|
private var index = 0
|
||||||
private var txStatusRequest: Events.TxStatusRequest? = null
|
private var txStatusRequest: Events.TxStatusRequest? = null
|
||||||
|
|
||||||
fun withRequest(req: BlockchainOuterClass.TxStatusRequest): TxStatus {
|
override fun onRequest(msg: BlockchainOuterClass.TxStatusRequest) {
|
||||||
this.txStatusRequest = Events.TxStatusRequest(req.txId)
|
this.txStatusRequest = Events.TxStatusRequest(msg.txId)
|
||||||
return withChain(req.chainValue)
|
withChain(msg.chainValue)
|
||||||
}
|
}
|
||||||
|
|
||||||
fun onReply(resp: BlockchainOuterClass.TxStatus): Events.TxStatus {
|
override fun onReply(msg: BlockchainOuterClass.TxStatus): Events.TxStatus {
|
||||||
return Events.TxStatus(
|
return Events.TxStatus(
|
||||||
chain, UUID.randomUUID(), requestDetails, txStatusRequest!!,
|
chain, UUID.randomUUID(), requestDetails, txStatusRequest!!,
|
||||||
Events.TxStatusResponse(resp.confirmations),
|
Events.TxStatusResponse(msg.confirmations),
|
||||||
index++
|
index++
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
@@ -183,7 +203,9 @@ class EventsBuilder {
|
|||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
class NativeCall : Base<NativeCall>() {
|
class NativeCall :
|
||||||
|
Base<NativeCall>(),
|
||||||
|
RequestReply<Events.NativeCall, BlockchainOuterClass.NativeCallRequest, BlockchainOuterClass.NativeCallReplyItem> {
|
||||||
val items = ArrayList<Events.NativeCallItemDetails>()
|
val items = ArrayList<Events.NativeCallItemDetails>()
|
||||||
val replies = HashMap<Int, Events.NativeCallReplyDetails>()
|
val replies = HashMap<Int, Events.NativeCallReplyDetails>()
|
||||||
private var index = 0
|
private var index = 0
|
||||||
@@ -192,24 +214,26 @@ class EventsBuilder {
|
|||||||
return this
|
return this
|
||||||
}
|
}
|
||||||
|
|
||||||
fun onRequest(item: BlockchainOuterClass.NativeCallItem): NativeCall {
|
override fun onRequest(msg: BlockchainOuterClass.NativeCallRequest) {
|
||||||
this.items.add(
|
withChain(msg.chain.number)
|
||||||
Events.NativeCallItemDetails(
|
msg.itemsList.forEach { item ->
|
||||||
item.method,
|
this.items.add(
|
||||||
item.id,
|
Events.NativeCallItemDetails(
|
||||||
item.payload.size().toLong()
|
item.method,
|
||||||
)
|
item.id,
|
||||||
)
|
item.payload.size().toLong()
|
||||||
return this
|
)
|
||||||
|
)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fun onReply(reply: BlockchainOuterClass.NativeCallReplyItem): Events.NativeCall {
|
override fun onReply(msg: BlockchainOuterClass.NativeCallReplyItem): Events.NativeCall {
|
||||||
val item = items.find { it.id == reply.id }!!
|
val item = items.find { it.id == msg.id }!!
|
||||||
return Events.NativeCall(
|
return Events.NativeCall(
|
||||||
request = requestDetails,
|
request = requestDetails,
|
||||||
total = items.size,
|
total = items.size,
|
||||||
index = index++,
|
index = index++,
|
||||||
succeed = reply.succeed,
|
succeed = msg.succeed,
|
||||||
blockchain = chain,
|
blockchain = chain,
|
||||||
nativeCall = item,
|
nativeCall = item,
|
||||||
payloadSizeBytes = item.payloadSizeBytes,
|
payloadSizeBytes = item.payloadSizeBytes,
|
||||||
@@ -219,13 +243,18 @@ class EventsBuilder {
|
|||||||
|
|
||||||
}
|
}
|
||||||
|
|
||||||
class Describe : Base<Describe>() {
|
class Describe :
|
||||||
|
Base<Describe>(),
|
||||||
|
RequestReply<Events.Describe, BlockchainOuterClass.DescribeRequest, BlockchainOuterClass.DescribeResponse> {
|
||||||
|
|
||||||
override fun getT(): Describe {
|
override fun getT(): Describe {
|
||||||
return this
|
return this
|
||||||
}
|
}
|
||||||
|
|
||||||
fun onReply(): Events.Describe {
|
override fun onRequest(msg: BlockchainOuterClass.DescribeRequest) {
|
||||||
|
}
|
||||||
|
|
||||||
|
override fun onReply(msg: BlockchainOuterClass.DescribeResponse): Events.Describe {
|
||||||
return Events.Describe(
|
return Events.Describe(
|
||||||
id = UUID.randomUUID(),
|
id = UUID.randomUUID(),
|
||||||
request = requestDetails
|
request = requestDetails
|
||||||
@@ -233,13 +262,18 @@ class EventsBuilder {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
class Status : Base<Status>() {
|
class Status :
|
||||||
|
Base<Status>(),
|
||||||
|
RequestReply<Events.Status, BlockchainOuterClass.StatusRequest, BlockchainOuterClass.ChainStatus> {
|
||||||
override fun getT(): Status {
|
override fun getT(): Status {
|
||||||
return this
|
return this
|
||||||
}
|
}
|
||||||
|
|
||||||
fun onReply(message: BlockchainOuterClass.ChainStatus): Events.Status {
|
override fun onRequest(msg: BlockchainOuterClass.StatusRequest) {
|
||||||
val chain = Chain.byId(message.chainValue)
|
}
|
||||||
|
|
||||||
|
override fun onReply(msg: BlockchainOuterClass.ChainStatus): Events.Status {
|
||||||
|
val chain = Chain.byId(msg.chainValue)
|
||||||
return Events.Status(
|
return Events.Status(
|
||||||
blockchain = chain,
|
blockchain = chain,
|
||||||
request = requestDetails,
|
request = requestDetails,
|
||||||
|
|||||||
@@ -20,10 +20,40 @@ import io.emeraldpay.grpc.Chain
|
|||||||
import io.grpc.Attributes
|
import io.grpc.Attributes
|
||||||
import io.grpc.Grpc
|
import io.grpc.Grpc
|
||||||
import io.grpc.Metadata
|
import io.grpc.Metadata
|
||||||
|
import org.jetbrains.annotations.NotNull
|
||||||
|
import org.junit.validator.TestClassValidator
|
||||||
import spock.lang.Specification
|
import spock.lang.Specification
|
||||||
|
|
||||||
class EventsBaseBuilderSpec extends Specification {
|
class EventsBaseBuilderSpec extends Specification {
|
||||||
|
|
||||||
|
class TestEvent extends Events.Base {
|
||||||
|
Events.StreamRequestDetails request
|
||||||
|
|
||||||
|
TestEvent(Events.StreamRequestDetails request) {
|
||||||
|
super(UUID.randomUUID(), "TEST")
|
||||||
|
this.request = request
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
class TestEventBuilder extends EventsBuilder.Base<TestEventBuilder>
|
||||||
|
implements EventsBuilder.RequestReply<TestEvent, BlockchainOuterClass.NativeCallRequest, BlockchainOuterClass.NativeCallReplyItem> {
|
||||||
|
|
||||||
|
@Override
|
||||||
|
protected TestEventBuilder getT() {
|
||||||
|
return this
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
void onRequest(BlockchainOuterClass.NativeCallRequest msg) {
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
TestEvent onReply(BlockchainOuterClass.NativeCallReplyItem msg) {
|
||||||
|
return new TestEvent(requestDetails)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
def "Parse headers from direct local access"() {
|
def "Parse headers from direct local access"() {
|
||||||
setup:
|
setup:
|
||||||
def metadata = new Metadata()
|
def metadata = new Metadata()
|
||||||
@@ -32,10 +62,11 @@ class EventsBaseBuilderSpec extends Specification {
|
|||||||
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("127.0.0.1"), 2448))
|
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("127.0.0.1"), 2448))
|
||||||
.build()
|
.build()
|
||||||
when:
|
when:
|
||||||
def act = new EventsBuilder.NativeCall()
|
def act = new TestEventBuilder()
|
||||||
.start(metadata, attributes)
|
.tap {
|
||||||
.withChain(Chain.ETHEREUM.id)
|
it.start(metadata, attributes)
|
||||||
.onRequest(BlockchainOuterClass.NativeCallItem.getDefaultInstance())
|
it.onRequest(BlockchainOuterClass.NativeCallRequest.getDefaultInstance())
|
||||||
|
}
|
||||||
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
||||||
then:
|
then:
|
||||||
act.request != null
|
act.request != null
|
||||||
@@ -56,10 +87,11 @@ class EventsBaseBuilderSpec extends Specification {
|
|||||||
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("127.0.0.1"), 2448))
|
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("127.0.0.1"), 2448))
|
||||||
.build()
|
.build()
|
||||||
when:
|
when:
|
||||||
def act = new EventsBuilder.NativeCall()
|
def act = new TestEventBuilder()
|
||||||
.start(metadata, attributes)
|
.tap {
|
||||||
.withChain(Chain.ETHEREUM.id)
|
it.start(metadata, attributes)
|
||||||
.onRequest(BlockchainOuterClass.NativeCallItem.getDefaultInstance())
|
it.onRequest(BlockchainOuterClass.NativeCallRequest.getDefaultInstance())
|
||||||
|
}
|
||||||
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
||||||
then:
|
then:
|
||||||
with(act.request.remote) {
|
with(act.request.remote) {
|
||||||
@@ -77,10 +109,11 @@ class EventsBaseBuilderSpec extends Specification {
|
|||||||
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448))
|
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448))
|
||||||
.build()
|
.build()
|
||||||
when:
|
when:
|
||||||
def act = new EventsBuilder.NativeCall()
|
def act = new TestEventBuilder()
|
||||||
.start(metadata, attributes)
|
.tap {
|
||||||
.withChain(Chain.ETHEREUM.id)
|
it.start(metadata, attributes)
|
||||||
.onRequest(BlockchainOuterClass.NativeCallItem.getDefaultInstance())
|
it.onRequest(BlockchainOuterClass.NativeCallRequest.getDefaultInstance())
|
||||||
|
}
|
||||||
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
||||||
then:
|
then:
|
||||||
with(act.request.remote) {
|
with(act.request.remote) {
|
||||||
@@ -99,10 +132,11 @@ class EventsBaseBuilderSpec extends Specification {
|
|||||||
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448))
|
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448))
|
||||||
.build()
|
.build()
|
||||||
when:
|
when:
|
||||||
def act = new EventsBuilder.NativeCall()
|
def act = new TestEventBuilder()
|
||||||
.start(metadata, attributes)
|
.tap {
|
||||||
.withChain(Chain.ETHEREUM.id)
|
it.start(metadata, attributes)
|
||||||
.onRequest(BlockchainOuterClass.NativeCallItem.getDefaultInstance())
|
it.onRequest(BlockchainOuterClass.NativeCallRequest.getDefaultInstance())
|
||||||
|
}
|
||||||
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
||||||
then:
|
then:
|
||||||
with(act.request.remote) {
|
with(act.request.remote) {
|
||||||
@@ -121,10 +155,11 @@ class EventsBaseBuilderSpec extends Specification {
|
|||||||
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448))
|
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448))
|
||||||
.build()
|
.build()
|
||||||
when:
|
when:
|
||||||
def act = new EventsBuilder.NativeCall()
|
def act = new TestEventBuilder()
|
||||||
.start(metadata, attributes)
|
.tap {
|
||||||
.withChain(Chain.ETHEREUM.id)
|
it.start(metadata, attributes)
|
||||||
.onRequest(BlockchainOuterClass.NativeCallItem.getDefaultInstance())
|
it.onRequest(BlockchainOuterClass.NativeCallRequest.getDefaultInstance())
|
||||||
|
}
|
||||||
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
||||||
then:
|
then:
|
||||||
with(act.request.remote) {
|
with(act.request.remote) {
|
||||||
@@ -143,10 +178,11 @@ class EventsBaseBuilderSpec extends Specification {
|
|||||||
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet6Address.getByName("::1"), 2448))
|
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet6Address.getByName("::1"), 2448))
|
||||||
.build()
|
.build()
|
||||||
when:
|
when:
|
||||||
def act = new EventsBuilder.NativeCall()
|
def act = new TestEventBuilder()
|
||||||
.start(metadata, attributes)
|
.tap {
|
||||||
.withChain(Chain.ETHEREUM.id)
|
it.start(metadata, attributes)
|
||||||
.onRequest(BlockchainOuterClass.NativeCallItem.getDefaultInstance())
|
it.onRequest(BlockchainOuterClass.NativeCallRequest.getDefaultInstance())
|
||||||
|
}
|
||||||
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
||||||
then:
|
then:
|
||||||
with(act.request.remote) {
|
with(act.request.remote) {
|
||||||
@@ -163,10 +199,11 @@ class EventsBaseBuilderSpec extends Specification {
|
|||||||
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448))
|
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448))
|
||||||
.build()
|
.build()
|
||||||
when:
|
when:
|
||||||
def act = new EventsBuilder.NativeCall()
|
def act = new TestEventBuilder()
|
||||||
.start(metadata, attributes)
|
.tap {
|
||||||
.withChain(Chain.ETHEREUM.id)
|
it.start(metadata, attributes)
|
||||||
.onRequest(BlockchainOuterClass.NativeCallItem.getDefaultInstance())
|
it.onRequest(BlockchainOuterClass.NativeCallRequest.getDefaultInstance())
|
||||||
|
}
|
||||||
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
||||||
then:
|
then:
|
||||||
with(act.request.remote) {
|
with(act.request.remote) {
|
||||||
@@ -183,10 +220,11 @@ class EventsBaseBuilderSpec extends Specification {
|
|||||||
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448))
|
.set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448))
|
||||||
.build()
|
.build()
|
||||||
when:
|
when:
|
||||||
def act = new EventsBuilder.NativeCall()
|
def act = new TestEventBuilder()
|
||||||
.start(metadata, attributes)
|
.tap {
|
||||||
.withChain(Chain.ETHEREUM.id)
|
it.start(metadata, attributes)
|
||||||
.onRequest(BlockchainOuterClass.NativeCallItem.getDefaultInstance())
|
it.onRequest(BlockchainOuterClass.NativeCallRequest.getDefaultInstance())
|
||||||
|
}
|
||||||
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
.onReply(BlockchainOuterClass.NativeCallReplyItem.getDefaultInstance())
|
||||||
then:
|
then:
|
||||||
with(act.request.remote) {
|
with(act.request.remote) {
|
||||||
|
|||||||
@@ -32,9 +32,9 @@ class EventsBuilderSubscribeBalanceSpec extends Specification {
|
|||||||
.setBalance("1234560000000000000")
|
.setBalance("1234560000000000000")
|
||||||
.build()
|
.build()
|
||||||
when:
|
when:
|
||||||
def act = new EventsBuilder.SubscribeBalance(true)
|
def act = new EventsBuilder.SubscribeBalance(true).tap {
|
||||||
.withRequest(request)
|
it.onRequest(request)
|
||||||
.onReply(resp)
|
}.onReply(resp)
|
||||||
then:
|
then:
|
||||||
act.index == 0
|
act.index == 0
|
||||||
act.blockchain == Chain.ETHEREUM
|
act.blockchain == Chain.ETHEREUM
|
||||||
@@ -69,9 +69,9 @@ class EventsBuilderSubscribeBalanceSpec extends Specification {
|
|||||||
.setBalance("12345600000000")
|
.setBalance("12345600000000")
|
||||||
.build()
|
.build()
|
||||||
when:
|
when:
|
||||||
def act = new EventsBuilder.SubscribeBalance(true)
|
def act = new EventsBuilder.SubscribeBalance(true).tap {
|
||||||
.withRequest(request)
|
it.onRequest(request)
|
||||||
.onReply(resp)
|
}.onReply(resp)
|
||||||
then:
|
then:
|
||||||
act.index == 0
|
act.index == 0
|
||||||
act.blockchain == Chain.BITCOIN
|
act.blockchain == Chain.BITCOIN
|
||||||
|
|||||||
Reference in New Issue
Block a user