From 0a3053d8034b5627301eb936ba4f3aa7b987d1c5 Mon Sep 17 00:00:00 2001 From: Igor Artamonov Date: Mon, 28 Jun 2021 16:47:41 -0400 Subject: [PATCH] solution: access logging for SubscribeBalance method --- .../monitoring/accesslog/AccessHandler.kt | 59 ++++++++++--- .../dshackle/monitoring/accesslog/Events.kt | 27 ++++-- .../monitoring/accesslog/EventsBuilder.kt | 29 ++++++- .../EventsBuilderSubscribeBalanceSpec.groovy | 83 +++++++++++++++++++ 4 files changed, 180 insertions(+), 18 deletions(-) create mode 100644 src/test/groovy/io/emeraldpay/dshackle/monitoring/accesslog/EventsBuilderSubscribeBalanceSpec.groovy diff --git a/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/AccessHandler.kt b/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/AccessHandler.kt index 5208dc89..e013727b 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/AccessHandler.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/AccessHandler.kt @@ -36,20 +36,15 @@ class AccessHandler( headers: Metadata, next: ServerCallHandler): ServerCall.Listener { - when (val method = call.methodDescriptor.bareMethodName) { - "SubscribeHead" -> { - return processSubscribeHead(call, headers, next) - } - "NativeCall" -> { - return processNativeCall(call, headers, next) - } + return when (val method = call.methodDescriptor.bareMethodName) { + "SubscribeHead" -> processSubscribeHead(call, headers, next) + "SubscribeBalance" -> processSubscribeBalance(call, headers, next) + "NativeCall" -> processNativeCall(call, headers, next) else -> { log.trace("unsupported method `{}`", method) + next.startCall(call, headers) } } - - // continue - return next.startCall(call, headers) } @Suppress("UNCHECKED_CAST") @@ -68,6 +63,22 @@ class AccessHandler( ) as ServerCall.Listener } + @Suppress("UNCHECKED_CAST") + private fun processSubscribeBalance( + call: ServerCall, + headers: Metadata, + next: ServerCallHandler + ): ServerCall.Listener { + val builder = EventsBuilder.SubscribeBalance() + .start(headers, call.attributes) + val callWrapper: ServerCall = OnSubscribeBalanceResponse( + call as ServerCall, builder, accessLogWriter) as ServerCall + return OnSubscribeBalance( + next.startCall(callWrapper, headers) as ServerCall.Listener, + builder + ) as ServerCall.Listener + } + @Suppress("UNCHECKED_CAST") private fun processNativeCall( call: ServerCall, @@ -103,6 +114,21 @@ class AccessHandler( } } + class OnSubscribeBalance( + val next: ServerCall.Listener, + val builder: EventsBuilder.SubscribeBalance + ) : ForwardingServerCallListener() { + + override fun onMessage(message: BlockchainOuterClass.BalanceRequest) { + builder.withRequest(message) + super.onMessage(message) + } + + override fun delegate(): ServerCall.Listener { + return next + } + } + class OnNativeCall( val next: ServerCall.Listener, val builder: EventsBuilder.NativeCall, @@ -172,4 +198,17 @@ class AccessHandler( super.sendMessage(message) } } + + class OnSubscribeBalanceResponse( + next: ServerCall, + val builder: EventsBuilder.SubscribeBalance, + val accessLogWriter: AccessLogWriter + ) : BaseCallResponse(next) { + + override fun sendMessage(message: BlockchainOuterClass.AddressBalance) { + val event = builder.onReply(message) + accessLogWriter.submit(event) + super.sendMessage(message) + } + } } \ No newline at end of file diff --git a/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/Events.kt b/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/Events.kt index 020453d1..13e211cb 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/Events.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/Events.kt @@ -16,15 +16,8 @@ package io.emeraldpay.dshackle.monitoring.accesslog import com.fasterxml.jackson.annotation.JsonInclude -import io.emeraldpay.api.proto.BlockchainOuterClass import io.emeraldpay.grpc.Chain -import io.grpc.Attributes -import io.grpc.Grpc -import io.grpc.Metadata -import org.apache.commons.lang3.StringUtils import org.slf4j.LoggerFactory -import java.net.InetAddress -import java.net.InetSocketAddress import java.time.Instant import java.util.* @@ -53,6 +46,17 @@ class Events { val index: Int ) : ChainBase(blockchain, "SubscribeHead", id) + @JsonInclude(JsonInclude.Include.NON_NULL) + class SubscribeBalance( + blockchain: Chain, id: UUID, + // initial request details + val request: StreamRequestDetails, + val balanceRequest: BalanceRequest, + val addressBalance: AddressBalance, + // index of the current response + val index: Int + ) : ChainBase(blockchain, "SubscribeBalance", id) + @JsonInclude(JsonInclude.Include.NON_NULL) class NativeCall( blockchain: Chain, id: UUID, @@ -98,4 +102,13 @@ class Events { val ts: Instant = Instant.now() ) + data class BalanceRequest( + val asset: String, + val addressType: String + ) + + data class AddressBalance( + val asset: String, + val address: String + ) } \ No newline at end of file diff --git a/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/EventsBuilder.kt b/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/EventsBuilder.kt index c9e223f3..093b5347 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/EventsBuilder.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/EventsBuilder.kt @@ -132,8 +132,35 @@ class EventsBuilder { } } - class NativeCall : Base() { + class SubscribeBalance() : Base() { + private var index = 0 + private var balanceRequest: Events.BalanceRequest? = null + override fun getT(): SubscribeBalance { + return this + } + + fun withRequest(req: BlockchainOuterClass.BalanceRequest): SubscribeBalance { + balanceRequest = Events.BalanceRequest( + req.asset.code.toUpperCase(), + req.address.addrTypeCase.name + ) + return this + } + + fun onReply(resp: BlockchainOuterClass.AddressBalance): Events.SubscribeBalance { + if (balanceRequest == null) { + throw IllegalStateException("Request is not initialized") + } + val addressBalance = Events.AddressBalance(resp.asset.code, resp.address.address) + val chain = Chain.byId(resp.asset.chain.number) + return Events.SubscribeBalance( + chain, UUID.randomUUID(), requestDetails, balanceRequest!!, addressBalance, index++ + ) + } + } + + class NativeCall : Base() { val items = ArrayList() val replies = HashMap() diff --git a/src/test/groovy/io/emeraldpay/dshackle/monitoring/accesslog/EventsBuilderSubscribeBalanceSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/monitoring/accesslog/EventsBuilderSubscribeBalanceSpec.groovy new file mode 100644 index 00000000..c445c590 --- /dev/null +++ b/src/test/groovy/io/emeraldpay/dshackle/monitoring/accesslog/EventsBuilderSubscribeBalanceSpec.groovy @@ -0,0 +1,83 @@ +package io.emeraldpay.dshackle.monitoring.accesslog + +import io.emeraldpay.api.proto.BlockchainOuterClass +import io.emeraldpay.api.proto.Common +import io.emeraldpay.grpc.Chain +import spock.lang.Specification + +class EventsBuilderSubscribeBalanceSpec extends Specification { + + def "Basic ethereum event"() { + setup: + def request = BlockchainOuterClass.BalanceRequest.newBuilder() + .setAddress( + Common.AnyAddress.newBuilder() + .setAddressSingle( + Common.SingleAddress.newBuilder() + .setAddress("0x7a250d5630B4cF539739dF2C5dAcb4c659F2488D") + ) + ) + .setAsset( + Common.Asset.newBuilder() + .setChainValue(100) + .setCode("ETHER") + ) + .build() + def resp = BlockchainOuterClass.AddressBalance.newBuilder() + .setAddress(Common.SingleAddress.newBuilder() + .setAddress("0x7a250d5630B4cF539739dF2C5dAcb4c659F2488D")) + .setAsset(Common.Asset.newBuilder() + .setChainValue(100) + .setCode("ETHER")) + .setBalance("1234560000000000000") + .build() + when: + def act = new EventsBuilder.SubscribeBalance() + .withRequest(request) + .onReply(resp) + then: + act.index == 0 + act.blockchain == Chain.ETHEREUM + act.balanceRequest.asset == "ETHER" + act.balanceRequest.addressType == "ADDRESS_SINGLE" + act.addressBalance.asset == "ETHER" + act.addressBalance.address == "0x7a250d5630B4cF539739dF2C5dAcb4c659F2488D" + } + + def "Basic bitcoin event"() { + setup: + def request = BlockchainOuterClass.BalanceRequest.newBuilder() + .setAddress( + Common.AnyAddress.newBuilder() + .setAddressSingle( + Common.SingleAddress.newBuilder() + .setAddress("1NDyJtNTjmwk5xPNhjgAMu4HDHigtobu1s") + ) + ) + .setAsset( + Common.Asset.newBuilder() + .setChainValue(1) + .setCode("BTC") + ) + .build() + def resp = BlockchainOuterClass.AddressBalance.newBuilder() + .setAddress(Common.SingleAddress.newBuilder() + .setAddress("1NDyJtNTjmwk5xPNhjgAMu4HDHigtobu1s")) + .setAsset(Common.Asset.newBuilder() + .setChainValue(1) + .setCode("BTC")) + .setBalance("12345600000000") + .build() + when: + def act = new EventsBuilder.SubscribeBalance() + .withRequest(request) + .onReply(resp) + then: + act.index == 0 + act.blockchain == Chain.BITCOIN + act.balanceRequest.asset == "BTC" + act.balanceRequest.addressType == "ADDRESS_SINGLE" + act.addressBalance.asset == "BTC" + act.addressBalance.address == "1NDyJtNTjmwk5xPNhjgAMu4HDHigtobu1s" + } +}