solution: access logging for SubscribeBalance method

This commit is contained in:
Igor Artamonov
2021-06-28 16:47:41 -04:00
parent ec043592c4
commit 0a3053d803
4 changed files with 180 additions and 18 deletions

View File

@@ -36,20 +36,15 @@ class AccessHandler(
headers: Metadata,
next: ServerCallHandler<ReqT, RespT>): ServerCall.Listener<ReqT> {
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<ReqT>
}
@Suppress("UNCHECKED_CAST")
private fun <ReqT : Any, RespT : Any> processSubscribeBalance(
call: ServerCall<ReqT, RespT>,
headers: Metadata,
next: ServerCallHandler<ReqT, RespT>
): ServerCall.Listener<ReqT> {
val builder = EventsBuilder.SubscribeBalance()
.start(headers, call.attributes)
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")
private fun <ReqT : Any, RespT : Any> processNativeCall(
call: ServerCall<ReqT, RespT>,
@@ -103,6 +114,21 @@ class AccessHandler(
}
}
class OnSubscribeBalance(
val next: ServerCall.Listener<BlockchainOuterClass.BalanceRequest>,
val builder: EventsBuilder.SubscribeBalance
) : ForwardingServerCallListener<BlockchainOuterClass.BalanceRequest>() {
override fun onMessage(message: BlockchainOuterClass.BalanceRequest) {
builder.withRequest(message)
super.onMessage(message)
}
override fun delegate(): ServerCall.Listener<BlockchainOuterClass.BalanceRequest> {
return next
}
}
class OnNativeCall(
val next: ServerCall.Listener<BlockchainOuterClass.NativeCallRequest>,
val builder: EventsBuilder.NativeCall,
@@ -172,4 +198,17 @@ class AccessHandler(
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)
}
}
}

View File

@@ -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
)
}

View File

@@ -132,8 +132,35 @@ class EventsBuilder {
}
}
class NativeCall : Base<NativeCall>() {
class SubscribeBalance() : Base<SubscribeBalance>() {
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<NativeCall>() {
val items = ArrayList<Events.NativeCallItemDetails>()
val replies = HashMap<Int, Events.NativeCallReplyDetails>()

View File

@@ -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"
}
}