solution: refactoring
This commit is contained in:
@@ -58,7 +58,7 @@ class AccessHandler(
|
||||
headers: Metadata,
|
||||
next: ServerCallHandler<ReqT, RespT>
|
||||
): ServerCall.Listener<ReqT> {
|
||||
val builder = Events.SubscribeHeadBuilder()
|
||||
val builder = EventsBuilder.SubscribeHead()
|
||||
.start(headers, call.attributes)
|
||||
val callWrapper: ServerCall<ReqT, RespT> = OnSubscribeHeadResponse(
|
||||
call as ServerCall<Common.Chain, BlockchainOuterClass.ChainHead>, builder, accessLogWriter) as ServerCall<ReqT, RespT>
|
||||
@@ -74,7 +74,7 @@ class AccessHandler(
|
||||
headers: Metadata,
|
||||
next: ServerCallHandler<ReqT, RespT>
|
||||
): ServerCall.Listener<ReqT> {
|
||||
val builder = Events.NativeCallBuilder()
|
||||
val builder = EventsBuilder.NativeCall()
|
||||
.start(headers, call.attributes)
|
||||
|
||||
val callWrapper: ServerCall<ReqT, RespT> = OnNativeCallResponse(
|
||||
@@ -89,7 +89,7 @@ class AccessHandler(
|
||||
|
||||
class OnSubscribeHead(
|
||||
val next: ServerCall.Listener<Common.Chain>,
|
||||
val builder: Events.SubscribeHeadBuilder
|
||||
val builder: EventsBuilder.SubscribeHead
|
||||
) : ForwardingServerCallListener<Common.Chain>() {
|
||||
|
||||
override fun onMessage(message: Common.Chain) {
|
||||
@@ -105,7 +105,7 @@ class AccessHandler(
|
||||
|
||||
class OnNativeCall(
|
||||
val next: ServerCall.Listener<BlockchainOuterClass.NativeCallRequest>,
|
||||
val builder: Events.NativeCallBuilder,
|
||||
val builder: EventsBuilder.NativeCall,
|
||||
val done: (List<Events.NativeCall>) -> Unit
|
||||
) : ForwardingServerCallListener<BlockchainOuterClass.NativeCallRequest>() {
|
||||
|
||||
@@ -151,7 +151,7 @@ class AccessHandler(
|
||||
|
||||
class OnNativeCallResponse(
|
||||
next: ServerCall<BlockchainOuterClass.NativeCallRequest, BlockchainOuterClass.NativeCallReplyItem>,
|
||||
val builder: Events.NativeCallBuilder
|
||||
val builder: EventsBuilder.NativeCall
|
||||
) : BaseCallResponse<BlockchainOuterClass.NativeCallRequest, BlockchainOuterClass.NativeCallReplyItem>(next) {
|
||||
|
||||
override fun sendMessage(message: BlockchainOuterClass.NativeCallReplyItem) {
|
||||
@@ -162,7 +162,7 @@ class AccessHandler(
|
||||
|
||||
class OnSubscribeHeadResponse(
|
||||
next: ServerCall<Common.Chain, BlockchainOuterClass.ChainHead>,
|
||||
val builder: Events.SubscribeHeadBuilder,
|
||||
val builder: EventsBuilder.SubscribeHead,
|
||||
val accessLogWriter: AccessLogWriter
|
||||
) : BaseCallResponse<Common.Chain, BlockchainOuterClass.ChainHead>(next) {
|
||||
|
||||
|
||||
@@ -98,148 +98,4 @@ class Events {
|
||||
val ts: Instant = Instant.now()
|
||||
)
|
||||
|
||||
abstract class BaseBuilder<T>() {
|
||||
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>): 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<InetAddress>()
|
||||
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<SubscribeHeadBuilder>() {
|
||||
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<NativeCallBuilder>() {
|
||||
|
||||
val items = ArrayList<NativeCallItemDetails>()
|
||||
val replies = HashMap<Int, NativeCallReplyDetails>()
|
||||
|
||||
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<NativeCall> {
|
||||
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()
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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<T>() {
|
||||
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>): 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<InetAddress>()
|
||||
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<SubscribeHead>() {
|
||||
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<NativeCall>() {
|
||||
|
||||
val items = ArrayList<Events.NativeCallItemDetails>()
|
||||
val replies = HashMap<Int, Events.NativeCallReplyDetails>()
|
||||
|
||||
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<Events.NativeCall> {
|
||||
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()
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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())
|
||||
|
||||
Reference in New Issue
Block a user