solution: access logging for SubscribeHead method

This commit is contained in:
Igor Artamonov
2021-06-27 18:44:41 -04:00
parent fd26d79f6e
commit a541c32104
3 changed files with 143 additions and 49 deletions

View File

@@ -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.grpc.* import io.grpc.*
import org.slf4j.LoggerFactory import org.slf4j.LoggerFactory
import org.springframework.beans.factory.annotation.Autowired import org.springframework.beans.factory.annotation.Autowired
@@ -36,14 +37,11 @@ class AccessHandler(
next: ServerCallHandler<ReqT, RespT>): ServerCall.Listener<ReqT> { next: ServerCallHandler<ReqT, RespT>): ServerCall.Listener<ReqT> {
when (val method = call.methodDescriptor.bareMethodName) { when (val method = call.methodDescriptor.bareMethodName) {
"SubscribeHead" -> {
return processSubscribeHead(call, headers, next)
}
"NativeCall" -> { "NativeCall" -> {
val builder = Events.NativeCallBuilder() return processNativeCall(call, headers, next)
.start(headers, call.attributes)
return OnNativeCall<ReqT, RespT>(
next.startCall(OnNativeCallResponse(call, builder), headers),
builder) { logs ->
accessLogWriter.submit(logs)
}
} }
else -> { else -> {
log.trace("unsupported method `{}`", method) log.trace("unsupported method `{}`", method)
@@ -54,20 +52,68 @@ class AccessHandler(
return next.startCall(call, headers) return next.startCall(call, headers)
} }
@Suppress("UNCHECKED_CAST")
private fun <ReqT : Any, RespT : Any> processSubscribeHead(
call: ServerCall<ReqT, RespT>,
headers: Metadata,
next: ServerCallHandler<ReqT, RespT>
): ServerCall.Listener<ReqT> {
val builder = Events.SubscribeHeadBuilder()
.start(headers, call.attributes)
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>
}
class OnNativeCall<ReqT : Any, RespT : Any>( @Suppress("UNCHECKED_CAST")
val next: ServerCall.Listener<ReqT>, private fun <ReqT : Any, RespT : Any> processNativeCall(
call: ServerCall<ReqT, RespT>,
headers: Metadata,
next: ServerCallHandler<ReqT, RespT>
): ServerCall.Listener<ReqT> {
val builder = Events.NativeCallBuilder()
.start(headers, call.attributes)
val callWrapper: ServerCall<ReqT, RespT> = OnNativeCallResponse(
call as ServerCall<BlockchainOuterClass.NativeCallRequest, BlockchainOuterClass.NativeCallReplyItem>, builder
) as ServerCall<ReqT, RespT>
return OnNativeCall(
next.startCall(callWrapper, headers) as ServerCall.Listener<BlockchainOuterClass.NativeCallRequest>,
builder) { logs ->
accessLogWriter.submit(logs)
} as ServerCall.Listener<ReqT>
}
class OnSubscribeHead(
val next: ServerCall.Listener<Common.Chain>,
val builder: Events.SubscribeHeadBuilder
) : ForwardingServerCallListener<Common.Chain>() {
override fun onMessage(message: Common.Chain) {
val chainId = message.type.number
builder.withChain(chainId)
super.onMessage(message)
}
override fun delegate(): ServerCall.Listener<Common.Chain> {
return next
}
}
class OnNativeCall(
val next: ServerCall.Listener<BlockchainOuterClass.NativeCallRequest>,
val builder: Events.NativeCallBuilder, val builder: Events.NativeCallBuilder,
val done: (List<Events.NativeCall>) -> Unit val done: (List<Events.NativeCall>) -> Unit
) : ForwardingServerCallListener<ReqT>() { ) : ForwardingServerCallListener<BlockchainOuterClass.NativeCallRequest>() {
override fun onMessage(message: ReqT) { override fun onMessage(message: BlockchainOuterClass.NativeCallRequest) {
if (message is BlockchainOuterClass.NativeCallRequest) { val chain = message.chain
val chain = message.chain builder.withChain(chain.number)
builder.withChain(chain.number) message.itemsList.forEach { item ->
message.itemsList.forEach { item -> builder.onItem(item)
builder.onItem(item)
}
} }
super.onMessage(message) super.onMessage(message)
} }
@@ -82,16 +128,14 @@ class AccessHandler(
done(builder.build()) done(builder.build())
} }
override fun delegate(): ServerCall.Listener<ReqT> { override fun delegate(): ServerCall.Listener<BlockchainOuterClass.NativeCallRequest> {
return next return next
} }
} }
class OnNativeCallResponse<ReqT : Any, RespT : Any>( abstract class BaseCallResponse<ReqT : Any, RespT : Any>(
val next: ServerCall<ReqT, RespT>, val next: ServerCall<ReqT, RespT>
val builder: Events.NativeCallBuilder
) : ForwardingServerCall<ReqT, RespT>() { ) : ForwardingServerCall<ReqT, RespT>() {
override fun getMethodDescriptor(): MethodDescriptor<ReqT, RespT> { override fun getMethodDescriptor(): MethodDescriptor<ReqT, RespT> {
return next.methodDescriptor return next.methodDescriptor
} }
@@ -101,9 +145,30 @@ class AccessHandler(
} }
override fun sendMessage(message: RespT) { override fun sendMessage(message: RespT) {
if (message is BlockchainOuterClass.NativeCallReplyItem) { super.sendMessage(message)
builder.onItemReply(message) }
} }
class OnNativeCallResponse(
next: ServerCall<BlockchainOuterClass.NativeCallRequest, BlockchainOuterClass.NativeCallReplyItem>,
val builder: Events.NativeCallBuilder
) : BaseCallResponse<BlockchainOuterClass.NativeCallRequest, BlockchainOuterClass.NativeCallReplyItem>(next) {
override fun sendMessage(message: BlockchainOuterClass.NativeCallReplyItem) {
builder.onItemReply(message)
super.sendMessage(message)
}
}
class OnSubscribeHeadResponse(
next: ServerCall<Common.Chain, BlockchainOuterClass.ChainHead>,
val builder: Events.SubscribeHeadBuilder,
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) super.sendMessage(message)
} }
} }

View File

@@ -35,20 +35,28 @@ class Events {
} }
abstract class Base( abstract class Base(
val method: String,
val id: UUID val id: UUID
) { ) {
val ts = Instant.now() val ts = Instant.now()
} }
abstract class ChainBase( abstract class ChainBase(
val blockchain: Chain, method: String, id: UUID val blockchain: Chain, val method: String, id: UUID
) : Base(method, id) { ) : Base(id)
} @JsonInclude(JsonInclude.Include.NON_NULL)
class SubscribeHead(
blockchain: Chain, id: UUID,
// initial request details
val request: StreamRequestDetails,
// index of the current response
val index: Int
) : ChainBase(blockchain, "SubscribeHead", id)
@JsonInclude(JsonInclude.Include.NON_NULL) @JsonInclude(JsonInclude.Include.NON_NULL)
class NativeCall( class NativeCall(
blockchain: Chain, id: UUID,
// info about the initial request, that may include several native calls // info about the initial request, that may include several native calls
val request: StreamRequestDetails, val request: StreamRequestDetails,
// total native calls passes within the initial request // total native calls passes within the initial request
@@ -62,11 +70,8 @@ class Events {
val succeed: Boolean, val succeed: Boolean,
val rpcError: Int? = null, val rpcError: Int? = null,
val payloadSizeBytes: Long, val payloadSizeBytes: Long,
val nativeCall: NativeCallItemDetails
blockchain: Chain, method: String, id: UUID ) : ChainBase(blockchain, "NativeCall", id)
) : ChainBase(blockchain, method, id) {
}
data class StreamRequestDetails( data class StreamRequestDetails(
val id: UUID, val id: UUID,
@@ -93,8 +98,7 @@ class Events {
val ts: Instant = Instant.now() val ts: Instant = Instant.now()
) )
class NativeCallBuilder() { abstract class BaseBuilder<T>() {
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),
@@ -103,15 +107,14 @@ class Events {
private val invalidCharacters = Regex("[\n\t]+") private val invalidCharacters = Regex("[\n\t]+")
} }
private var requestDetails = StreamRequestDetails( var requestDetails = StreamRequestDetails(
UUID.randomUUID(), UUID.randomUUID(),
Instant.now(), Instant.now(),
Remote(emptyList(), "", "") Remote(emptyList(), "", "")
) )
var chain: Int = Chain.UNSPECIFIED.id var chainId: Int = Chain.UNSPECIFIED.id
val items = ArrayList<NativeCallItemDetails>() var chain = Chain.UNSPECIFIED
val replies = HashMap<Int, NativeCallReplyDetails>()
private fun toInetAddress(ip: String): InetAddress? { private fun toInetAddress(ip: String): InetAddress? {
val isIp = Character.digit(ip[0], 16) != -1 val isIp = Character.digit(ip[0], 16) != -1
@@ -144,15 +147,17 @@ class Events {
.trim() .trim()
} }
fun start(metadata: Metadata, attributes: Attributes): NativeCallBuilder { abstract protected fun getT(): T
fun start(metadata: Metadata, attributes: Attributes): T {
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@NativeCallBuilder::clean) ?.let(this@BaseBuilder::clean)
?: "" ?: ""
val ips = ArrayList<InetAddress>() val ips = ArrayList<InetAddress>()
remoteIpKeys.forEach { key -> remoteIpKeys.forEach { key ->
metadata.get(key)?.let { metadata.get(key)?.let {
it.trim().ifEmpty { null } it.trim().ifEmpty { null }
?.let(this@NativeCallBuilder::toInetAddress) ?.let(this@BaseBuilder::toInetAddress)
?.let(ips::add) ?.let(ips::add)
} }
} }
@@ -168,11 +173,36 @@ class Events {
ip = ip, ip = ip,
userAgent = userAgent userAgent = userAgent
)) ))
return getT()
}
fun withChain(chain: Int): T {
this.chainId = chain
this.chain = Chain.byId(chainId)
return getT()
}
}
class SubscribeHeadBuilder() : BaseBuilder<SubscribeHeadBuilder>() {
private var index = 0
override fun getT(): SubscribeHeadBuilder {
return this return this
} }
fun withChain(chain: Int): NativeCallBuilder { fun onReply(resp: BlockchainOuterClass.ChainHead): SubscribeHead {
this.chain = chain return SubscribeHead(
chain, UUID.randomUUID(), requestDetails, index++
)
}
}
class NativeCallBuilder : BaseBuilder<NativeCallBuilder>() {
val items = ArrayList<NativeCallItemDetails>()
val replies = HashMap<Int, NativeCallReplyDetails>()
override fun getT(): NativeCallBuilder {
return this return this
} }
@@ -197,7 +227,6 @@ class Events {
} }
fun build(): List<NativeCall> { fun build(): List<NativeCall> {
val blockchain = Chain.byId(this.chain)
return items.mapIndexed { index, item -> return items.mapIndexed { index, item ->
val reply = replies[item.id] val reply = replies[item.id]
NativeCall( NativeCall(
@@ -205,8 +234,8 @@ class Events {
total = items.size, total = items.size,
index = index, index = index,
succeed = reply?.succeed ?: false, succeed = reply?.succeed ?: false,
blockchain = blockchain, blockchain = chain,
method = item.method, nativeCall = item,
payloadSizeBytes = item.payloadSizeBytes, payloadSizeBytes = item.payloadSizeBytes,
id = UUID.randomUUID() id = UUID.randomUUID()
) )

View File

@@ -22,7 +22,7 @@ import io.grpc.Grpc
import io.grpc.Metadata import io.grpc.Metadata
import spock.lang.Specification import spock.lang.Specification
class EventsNativeCallBuilderSpec extends Specification { class EventsBaseBuilderSpec extends Specification {
def "Parse headers from direct local access"() { def "Parse headers from direct local access"() {
setup: setup: