From ec043592c48e61705d5a1894cab11375eb2049be Mon Sep 17 00:00:00 2001 From: Igor Artamonov Date: Mon, 28 Jun 2021 13:58:34 -0400 Subject: [PATCH] solution: refactoring --- .../monitoring/accesslog/AccessHandler.kt | 12 +- .../dshackle/monitoring/accesslog/Events.kt | 144 -------------- .../monitoring/accesslog/EventsBuilder.kt | 180 ++++++++++++++++++ .../accesslog/EventsBaseBuilderSpec.groovy | 16 +- 4 files changed, 194 insertions(+), 158 deletions(-) create mode 100644 src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/EventsBuilder.kt 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 fd0ae1b9..5208dc89 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/AccessHandler.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/AccessHandler.kt @@ -58,7 +58,7 @@ class AccessHandler( headers: Metadata, next: ServerCallHandler ): ServerCall.Listener { - val builder = Events.SubscribeHeadBuilder() + val builder = EventsBuilder.SubscribeHead() .start(headers, call.attributes) val callWrapper: ServerCall = OnSubscribeHeadResponse( call as ServerCall, builder, accessLogWriter) as ServerCall @@ -74,7 +74,7 @@ class AccessHandler( headers: Metadata, next: ServerCallHandler ): ServerCall.Listener { - val builder = Events.NativeCallBuilder() + val builder = EventsBuilder.NativeCall() .start(headers, call.attributes) val callWrapper: ServerCall = OnNativeCallResponse( @@ -89,7 +89,7 @@ class AccessHandler( class OnSubscribeHead( val next: ServerCall.Listener, - val builder: Events.SubscribeHeadBuilder + val builder: EventsBuilder.SubscribeHead ) : ForwardingServerCallListener() { override fun onMessage(message: Common.Chain) { @@ -105,7 +105,7 @@ class AccessHandler( class OnNativeCall( val next: ServerCall.Listener, - val builder: Events.NativeCallBuilder, + val builder: EventsBuilder.NativeCall, val done: (List) -> Unit ) : ForwardingServerCallListener() { @@ -151,7 +151,7 @@ class AccessHandler( class OnNativeCallResponse( next: ServerCall, - val builder: Events.NativeCallBuilder + val builder: EventsBuilder.NativeCall ) : BaseCallResponse(next) { override fun sendMessage(message: BlockchainOuterClass.NativeCallReplyItem) { @@ -162,7 +162,7 @@ class AccessHandler( class OnSubscribeHeadResponse( next: ServerCall, - val builder: Events.SubscribeHeadBuilder, + val builder: EventsBuilder.SubscribeHead, val accessLogWriter: AccessLogWriter ) : BaseCallResponse(next) { 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 052b7aa5..020453d1 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/Events.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/Events.kt @@ -98,148 +98,4 @@ class Events { val ts: Instant = Instant.now() ) - abstract class BaseBuilder() { - companion object { - private val remoteIpKeys = listOf( - Metadata.Key.of("x-real-ip", Metadata.ASCII_STRING_MARSHALLER), - Metadata.Key.of("x-forwarded-for", Metadata.ASCII_STRING_MARSHALLER) - ) - private val invalidCharacters = Regex("[\n\t]+") - } - - var requestDetails = StreamRequestDetails( - UUID.randomUUID(), - Instant.now(), - Remote(emptyList(), "", "") - ) - - var chainId: Int = Chain.UNSPECIFIED.id - var chain = Chain.UNSPECIFIED - - private fun toInetAddress(ip: String): InetAddress? { - val isIp = Character.digit(ip[0], 16) != -1 - if (!isIp) { - return null - } - return try { - InetAddress.getByName(ip) - } catch (t: Throwable) { - null - } - } - - private fun findBestIp(ips: List): InetAddress? { - // check if a real remote address is provided, otherwise use any local address - return ips.sortedWith(kotlin.Comparator { a, b -> - val aLocal = a.isLoopbackAddress || a.isSiteLocalAddress - val bLocal = b.isLoopbackAddress || b.isSiteLocalAddress - when { - aLocal && bLocal -> 0 - aLocal -> 1 - else -> -1 - } - }).firstOrNull() - } - - private fun clean(s: String): String { - return StringUtils.truncate(s, 128) - .replace(invalidCharacters, " ") - .trim() - } - - 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)) - ?.let(this@BaseBuilder::clean) - ?: "" - val ips = ArrayList() - remoteIpKeys.forEach { key -> - metadata.get(key)?.let { - it.trim().ifEmpty { null } - ?.let(this@BaseBuilder::toInetAddress) - ?.let(ips::add) - } - } - attributes.get(Grpc.TRANSPORT_ATTR_REMOTE_ADDR)?.let { addr -> - if (addr is InetSocketAddress) { - ips.add(addr.address) - } - } - val ip = findBestIp(ips)?.hostAddress ?: "" - this.requestDetails = this.requestDetails - .copy(remote = Remote( - ips = ips.map { it.hostAddress }, - ip = ip, - userAgent = userAgent - )) - return getT() - } - - fun withChain(chain: Int): T { - this.chainId = chain - this.chain = Chain.byId(chainId) - return getT() - } - } - - class SubscribeHeadBuilder() : BaseBuilder() { - private var index = 0 - - override fun getT(): SubscribeHeadBuilder { - return this - } - - fun onReply(resp: BlockchainOuterClass.ChainHead): SubscribeHead { - return SubscribeHead( - chain, UUID.randomUUID(), requestDetails, index++ - ) - } - } - - class NativeCallBuilder : BaseBuilder() { - - val items = ArrayList() - val replies = HashMap() - - override fun getT(): NativeCallBuilder { - return this - } - - fun onItem(item: BlockchainOuterClass.NativeCallItem): NativeCallBuilder { - this.items.add( - NativeCallItemDetails( - item.method, - item.id, - item.payload.size().toLong() - ) - ) - return this - } - - fun onItemReply(reply: BlockchainOuterClass.NativeCallReplyItem): NativeCallBuilder { - this.replies[reply.id] = NativeCallReplyDetails( - reply.id, - reply.succeed, - reply.payload?.size()?.toLong() ?: 0L - ) - return this - } - - fun build(): List { - return items.mapIndexed { index, item -> - val reply = replies[item.id] - NativeCall( - request = requestDetails, - total = items.size, - index = index, - succeed = reply?.succeed ?: false, - blockchain = chain, - nativeCall = item, - payloadSizeBytes = item.payloadSizeBytes, - id = UUID.randomUUID() - ) - } - } - } } \ 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 new file mode 100644 index 00000000..c9e223f3 --- /dev/null +++ b/src/main/kotlin/io/emeraldpay/dshackle/monitoring/accesslog/EventsBuilder.kt @@ -0,0 +1,180 @@ +/** + * Copyright (c) 2021 EmeraldPay, Inc + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.emeraldpay.dshackle.monitoring.accesslog + +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.* + +class EventsBuilder { + + companion object { + private val log = LoggerFactory.getLogger(EventsBuilder::class.java) + } + + abstract class Base() { + companion object { + private val remoteIpKeys = listOf( + Metadata.Key.of("x-real-ip", Metadata.ASCII_STRING_MARSHALLER), + Metadata.Key.of("x-forwarded-for", Metadata.ASCII_STRING_MARSHALLER) + ) + private val invalidCharacters = Regex("[\n\t]+") + } + + var requestDetails = Events.StreamRequestDetails( + UUID.randomUUID(), + Instant.now(), + Events.Remote(emptyList(), "", "") + ) + + var chainId: Int = Chain.UNSPECIFIED.id + var chain = Chain.UNSPECIFIED + + private fun toInetAddress(ip: String): InetAddress? { + val isIp = Character.digit(ip[0], 16) != -1 + if (!isIp) { + return null + } + return try { + InetAddress.getByName(ip) + } catch (t: Throwable) { + null + } + } + + private fun findBestIp(ips: List): InetAddress? { + // check if a real remote address is provided, otherwise use any local address + return ips.sortedWith(kotlin.Comparator { a, b -> + val aLocal = a.isLoopbackAddress || a.isSiteLocalAddress + val bLocal = b.isLoopbackAddress || b.isSiteLocalAddress + when { + aLocal && bLocal -> 0 + aLocal -> 1 + else -> -1 + } + }).firstOrNull() + } + + private fun clean(s: String): String { + return StringUtils.truncate(s, 128) + .replace(invalidCharacters, " ") + .trim() + } + + 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)) + ?.let(this@Base::clean) + ?: "" + val ips = ArrayList() + remoteIpKeys.forEach { key -> + metadata.get(key)?.let { + it.trim().ifEmpty { null } + ?.let(this@Base::toInetAddress) + ?.let(ips::add) + } + } + attributes.get(Grpc.TRANSPORT_ATTR_REMOTE_ADDR)?.let { addr -> + if (addr is InetSocketAddress) { + ips.add(addr.address) + } + } + val ip = findBestIp(ips)?.hostAddress ?: "" + this.requestDetails = this.requestDetails + .copy(remote = Events.Remote( + ips = ips.map { it.hostAddress }, + ip = ip, + userAgent = userAgent + )) + return getT() + } + + fun withChain(chain: Int): T { + this.chainId = chain + this.chain = Chain.byId(chainId) + return getT() + } + } + + class SubscribeHead() : Base() { + private var index = 0 + + override fun getT(): SubscribeHead { + return this + } + + fun onReply(resp: BlockchainOuterClass.ChainHead): Events.SubscribeHead { + return Events.SubscribeHead( + chain, UUID.randomUUID(), requestDetails, index++ + ) + } + } + + class NativeCall : Base() { + + val items = ArrayList() + val replies = HashMap() + + override fun getT(): NativeCall { + return this + } + + fun onItem(item: BlockchainOuterClass.NativeCallItem): NativeCall { + this.items.add( + Events.NativeCallItemDetails( + item.method, + item.id, + item.payload.size().toLong() + ) + ) + return this + } + + fun onItemReply(reply: BlockchainOuterClass.NativeCallReplyItem): NativeCall { + this.replies[reply.id] = Events.NativeCallReplyDetails( + reply.id, + reply.succeed, + reply.payload?.size()?.toLong() ?: 0L + ) + return this + } + + fun build(): List { + return items.mapIndexed { index, item -> + val reply = replies[item.id] + Events.NativeCall( + request = requestDetails, + total = items.size, + index = index, + succeed = reply?.succeed ?: false, + blockchain = chain, + nativeCall = item, + payloadSizeBytes = item.payloadSizeBytes, + id = UUID.randomUUID() + ) + } + } + } +} \ No newline at end of file diff --git a/src/test/groovy/io/emeraldpay/dshackle/monitoring/accesslog/EventsBaseBuilderSpec.groovy b/src/test/groovy/io/emeraldpay/dshackle/monitoring/accesslog/EventsBaseBuilderSpec.groovy index 85a55244..b72ec97e 100644 --- a/src/test/groovy/io/emeraldpay/dshackle/monitoring/accesslog/EventsBaseBuilderSpec.groovy +++ b/src/test/groovy/io/emeraldpay/dshackle/monitoring/accesslog/EventsBaseBuilderSpec.groovy @@ -32,7 +32,7 @@ class EventsBaseBuilderSpec extends Specification { .set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("127.0.0.1"), 2448)) .build() when: - def act = new Events.NativeCallBuilder() + def act = new EventsBuilder.NativeCall() .start(metadata, attributes) .withChain(Chain.ETHEREUM.id) .onItem(BlockchainOuterClass.NativeCallItem.getDefaultInstance()) @@ -59,7 +59,7 @@ class EventsBaseBuilderSpec extends Specification { .set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("127.0.0.1"), 2448)) .build() when: - def act = new Events.NativeCallBuilder() + def act = new EventsBuilder.NativeCall() .start(metadata, attributes) .withChain(Chain.ETHEREUM.id) .onItem(BlockchainOuterClass.NativeCallItem.getDefaultInstance()) @@ -81,7 +81,7 @@ class EventsBaseBuilderSpec extends Specification { .set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448)) .build() when: - def act = new Events.NativeCallBuilder() + def act = new EventsBuilder.NativeCall() .start(metadata, attributes) .withChain(Chain.ETHEREUM.id) .onItem(BlockchainOuterClass.NativeCallItem.getDefaultInstance()) @@ -104,7 +104,7 @@ class EventsBaseBuilderSpec extends Specification { .set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448)) .build() when: - def act = new Events.NativeCallBuilder() + def act = new EventsBuilder.NativeCall() .start(metadata, attributes) .withChain(Chain.ETHEREUM.id) .onItem(BlockchainOuterClass.NativeCallItem.getDefaultInstance()) @@ -127,7 +127,7 @@ class EventsBaseBuilderSpec extends Specification { .set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448)) .build() when: - def act = new Events.NativeCallBuilder() + def act = new EventsBuilder.NativeCall() .start(metadata, attributes) .withChain(Chain.ETHEREUM.id) .onItem(BlockchainOuterClass.NativeCallItem.getDefaultInstance()) @@ -150,7 +150,7 @@ class EventsBaseBuilderSpec extends Specification { .set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet6Address.getByName("::1"), 2448)) .build() when: - def act = new Events.NativeCallBuilder() + def act = new EventsBuilder.NativeCall() .start(metadata, attributes) .withChain(Chain.ETHEREUM.id) .onItem(BlockchainOuterClass.NativeCallItem.getDefaultInstance()) @@ -171,7 +171,7 @@ class EventsBaseBuilderSpec extends Specification { .set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448)) .build() when: - def act = new Events.NativeCallBuilder() + def act = new EventsBuilder.NativeCall() .start(metadata, attributes) .withChain(Chain.ETHEREUM.id) .onItem(BlockchainOuterClass.NativeCallItem.getDefaultInstance()) @@ -192,7 +192,7 @@ class EventsBaseBuilderSpec extends Specification { .set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, new InetSocketAddress(Inet4Address.getByName("30.56.100.15"), 2448)) .build() when: - def act = new Events.NativeCallBuilder() + def act = new EventsBuilder.NativeCall() .start(metadata, attributes) .withChain(Chain.ETHEREUM.id) .onItem(BlockchainOuterClass.NativeCallItem.getDefaultInstance())