Skip to content
Merged

Dev #715

Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ spring:
instance-id: ${spring.application.name}:${server.port}
healthCheckInterval: 20s
prefer-ip-address: true
query-passing: true
config:
import: vault://secret/${spring.application.name}
management:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import co.nilin.opex.matching.engine.core.model.MatchConstraint
import co.nilin.opex.matching.engine.core.model.OrderDirection
import co.nilin.opex.matching.engine.core.model.OrderType
import java.math.BigDecimal
import java.time.LocalDateTime

data class RichOrder(
val orderId: Long? = 0,
Expand All @@ -23,5 +24,6 @@ data class RichOrder(
val quoteQuantity: BigDecimal,
val executedQuantity: BigDecimal,
val accumulativeQuoteQty: BigDecimal,
val status: Int = 0
val status: Int = 0,
val createDate: LocalDateTime?
) : RichOrderEvent
Original file line number Diff line number Diff line change
@@ -1,13 +1,15 @@
package co.nilin.opex.accountant.core.inout

import java.math.BigDecimal
import java.time.LocalDateTime

data class RichOrderUpdate(
val ouid: String,
val price: BigDecimal,
val quantity: BigDecimal,
val remainedQuantity: BigDecimal,
val status: OrderStatus = OrderStatus.NEW
val status: OrderStatus = OrderStatus.NEW,
val updateDate: LocalDateTime?= LocalDateTime.now()
) : RichOrderEvent {

fun executedQuantity(): BigDecimal = quantity.minus(remainedQuantity)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -212,12 +212,13 @@ open class OrderManagerImpl(
richOrderPublisher.publish(
RichOrderUpdate(
order.ouid,
order.price.toBigDecimal(),
order.quantity.toBigDecimal(),
cancelOrderEvent.remainedQuantity.toBigDecimal(),
order.price.toBigDecimal().multiply(order.rightSideFraction),
order.origQuantity,
cancelOrderEvent.remainedQuantity.toBigDecimal().multiply(order.leftSideFraction),
OrderStatus.CANCELED
)
)

return financialActionPersister.persist(listOf(financialAction))
/*publishFinancialAction(financialAction)
return fa*/
Expand Down Expand Up @@ -253,7 +254,8 @@ open class OrderManagerImpl(
OrderStatus.NEW.code
} else {
OrderStatus.PARTIALLY_FILLED.code
}
},
LocalDateTime.now()
)
)
}
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package co.nilin.opex.accountant.ports.kafka.listener.config
package co.nilin.opex.accountant.ports.kafka.listener.config

import co.nilin.opex.accountant.core.inout.KycLevelUpdatedEvent
import co.nilin.opex.accountant.ports.kafka.listener.consumer.*
Expand All @@ -9,7 +9,6 @@ import co.nilin.opex.matching.engine.core.eventh.events.CoreEvent
import org.apache.kafka.clients.consumer.ConsumerConfig
import org.apache.kafka.common.TopicPartition
import org.apache.kafka.common.serialization.StringDeserializer
import org.springframework.beans.factory.annotation.Autowired
import org.springframework.beans.factory.annotation.Qualifier
import org.springframework.beans.factory.annotation.Value
import org.springframework.boot.autoconfigure.condition.ConditionalOnBean
Expand Down Expand Up @@ -61,85 +60,83 @@ class AccountantKafkaConfig {
fun withdrawRequestConsumerFactory(@Qualifier("consumerConfig") consumerConfigs: Map<String, Any?>): ConsumerFactory<String, WithdrawRequestEvent> {
return DefaultKafkaConsumerFactory(consumerConfigs)
}

@Bean("depositConsumerFactory")
fun depositConsumerFactory(@Qualifier("consumerConfig") consumerConfigs: Map<String, Any?>): ConsumerFactory<String, DepositEvent> {
return DefaultKafkaConsumerFactory(consumerConfigs)
}

@Autowired
@Bean("tradeKafkaListenerContainer")
@ConditionalOnBean(TradeKafkaListener::class)
fun configureTradeListener(
fun tradeListenerContainer(
tradeListener: TradeKafkaListener,
@Qualifier("accountantEventKafkaTemplate") template: KafkaTemplate<String?, CoreEvent>,
@Qualifier("accountantConsumerFactory") consumerFactory: ConsumerFactory<String, CoreEvent>
) {
): ConcurrentMessageListenerContainer<String, CoreEvent> {
val containerProps = ContainerProperties(Pattern.compile("trades_.*"))
containerProps.messageListener = tradeListener
val container = ConcurrentMessageListenerContainer(consumerFactory, containerProps)
container.setBeanName("TradeKafkaListenerContainer")
container.commonErrorHandler = createConsumerErrorHandler(template, "trades.DLT")
container.start()
return container
}

@Autowired
@Bean("eventKafkaListenerContainer")
@ConditionalOnBean(EventKafkaListener::class)
fun configureEventListener(
fun eventListenerContainer(
eventListener: EventKafkaListener,
@Qualifier("accountantEventKafkaTemplate") template: KafkaTemplate<String?, CoreEvent>,
@Qualifier("accountantConsumerFactory") consumerFactory: ConsumerFactory<String, CoreEvent>
) {
): ConcurrentMessageListenerContainer<String, CoreEvent> {
val containerProps = ContainerProperties(Pattern.compile("events_.*"))
containerProps.messageListener = eventListener
val container = ConcurrentMessageListenerContainer(consumerFactory, containerProps)
container.setBeanName("EventKafkaListenerContainer")
container.commonErrorHandler = createConsumerErrorHandler(template, "events.DLT")
container.start()
return container
}

@Autowired
@Bean("orderKafkaListenerContainer")
@ConditionalOnBean(OrderKafkaListener::class)
fun configureOrderListener(
fun orderListenerContainer(
orderListener: OrderKafkaListener,
@Qualifier("accountantEventKafkaTemplate") template: KafkaTemplate<String?, CoreEvent>,
@Qualifier("accountantConsumerFactory") consumerFactory: ConsumerFactory<String, CoreEvent>
) {
): ConcurrentMessageListenerContainer<String, CoreEvent> {
val containerProps = ContainerProperties(Pattern.compile("orders_.*"))
containerProps.messageListener = orderListener
val container = ConcurrentMessageListenerContainer(consumerFactory, containerProps)
container.setBeanName("OrderKafkaListenerContainer")
container.commonErrorHandler = createConsumerErrorHandler(template, "orders.DLT")
container.start()
return container
}

@Autowired
@Bean("tempEventKafkaListenerContainer")
@ConditionalOnBean(TempEventKafkaListener::class)
fun configureTempEventListener(
fun tempEventListenerContainer(
eventListener: TempEventKafkaListener,
@Qualifier("accountantEventKafkaTemplate") template: KafkaTemplate<String?, CoreEvent>,
@Qualifier("accountantConsumerFactory") consumerFactory: ConsumerFactory<String, CoreEvent>
) {
): ConcurrentMessageListenerContainer<String, CoreEvent> {
val containerProps = ContainerProperties(Pattern.compile("tempevents"))
containerProps.messageListener = eventListener
val container = ConcurrentMessageListenerContainer(consumerFactory, containerProps)
container.setBeanName("TempEventKafkaListenerContainer")
container.commonErrorHandler = createConsumerErrorHandler(template, "tempevents.DLT")
container.start()
return container
}

@Autowired
@Bean("faResponseKafkaListenerContainer")
@ConditionalOnBean(FAResponseKafkaListener::class)
fun configureEventListener(
fun faResponseListenerContainer(
eventListener: FAResponseKafkaListener,
//@Qualifier("accountantEventKafkaTemplate") template: KafkaTemplate<String?, CoreEvent>,
@Qualifier("faResponseConsumerFactory") consumerFactory: ConsumerFactory<String, FinancialActionResponseEvent>
) {
): ConcurrentMessageListenerContainer<String, FinancialActionResponseEvent> {
val containerProps = ContainerProperties(Pattern.compile("fiAction_response"))
containerProps.messageListener = eventListener
val container = ConcurrentMessageListenerContainer(consumerFactory, containerProps)
container.setBeanName("FAResponseKafkaListenerContainer")
//TODO add error handler
//container.commonErrorHandler = createConsumerErrorHandler(template, "events.DLT")
container.start()
return container
}

@Bean("kycLevelUpdatedProducerFactory")
Expand All @@ -152,69 +149,69 @@ class AccountantKafkaConfig {
return KafkaTemplate(producerFactory)
}

@Bean("withdrawRequestProducerFactory")
fun withdrawRequestProducerFactory(@Qualifier("consumerConfig") producerConfigs: Map<String, Any>): ProducerFactory<String, WithdrawRequestEvent> {
return DefaultKafkaProducerFactory(producerConfigs)
}

@Bean("withdrawRequestKafkaTemplate")
fun withdrawRequestKafkaTemplate(@Qualifier("withdrawRequestProducerFactory") producerFactory: ProducerFactory<String, WithdrawRequestEvent>): KafkaTemplate<String, WithdrawRequestEvent> {
return KafkaTemplate(producerFactory)
}

@Bean("depositProducerFactory")
fun depositProducerFactory(@Qualifier("consumerConfig") producerConfigs: Map<String, Any>): ProducerFactory<String, DepositEvent> {
return DefaultKafkaProducerFactory(producerConfigs)
}

@Bean("depositKafkaTemplate")
fun depositKafkaTemplate(@Qualifier("depositProducerFactory") producerFactory: ProducerFactory<String, DepositEvent>): KafkaTemplate<String, DepositEvent> {
return KafkaTemplate(producerFactory)
}

@Autowired
@Bean("kycLevelUpdatedKafkaListenerContainer")
@ConditionalOnBean(KycLevelUpdatedKafkaListener::class)
fun configureKycLevelUpdatedListener(
fun kycListenerContainer(
listener: KycLevelUpdatedKafkaListener,
@Qualifier("kycLevelUpdatedKafkaTemplate") template: KafkaTemplate<String, KycLevelUpdatedEvent>,
@Qualifier("KycConsumerFactory") consumerFactory: ConsumerFactory<String, KycLevelUpdatedEvent>
) {
): ConcurrentMessageListenerContainer<String, KycLevelUpdatedEvent> {
val containerProps = ContainerProperties(Pattern.compile("kyc_level_updated"))
containerProps.messageListener = listener
val container = ConcurrentMessageListenerContainer(consumerFactory, containerProps)
container.setBeanName("KycLevelUpdatedKafkaListenerContainer")
container.commonErrorHandler = createConsumerErrorHandler(template, "kyc_level_updated.DLT")
container.start()
return container
}

@Autowired
@Bean("withdrawRequestProducerFactory")
fun withdrawRequestProducerFactory(@Qualifier("consumerConfig") producerConfigs: Map<String, Any>): ProducerFactory<String, WithdrawRequestEvent> {
return DefaultKafkaProducerFactory(producerConfigs)
}

@Bean("withdrawRequestKafkaTemplate")
fun withdrawRequestKafkaTemplate(@Qualifier("withdrawRequestProducerFactory") producerFactory: ProducerFactory<String, WithdrawRequestEvent>): KafkaTemplate<String, WithdrawRequestEvent> {
return KafkaTemplate(producerFactory)
}

@Bean("withdrawRequestKafkaListenerContainer")
@ConditionalOnBean(WithdrawRequestKafkaListener::class)
fun configureWithdrawRequestEventListener(
fun withdrawRequestListenerContainer(
listener: WithdrawRequestKafkaListener,
@Qualifier("withdrawRequestKafkaTemplate") template: KafkaTemplate<String, WithdrawRequestEvent>,
@Qualifier("withdrawRequestConsumerFactory") consumerFactory: ConsumerFactory<String, WithdrawRequestEvent>
) {
): ConcurrentMessageListenerContainer<String, WithdrawRequestEvent> {
val containerProps = ContainerProperties(Pattern.compile("withdraw_request"))
containerProps.messageListener = listener
val container = ConcurrentMessageListenerContainer(consumerFactory, containerProps)
container.setBeanName("WithdrawRequestKafkaListenerContainer")
container.commonErrorHandler = createConsumerErrorHandler(template, "withdraw_request.DLT")
container.start()
return container
}

@Autowired
@Bean("depositProducerFactory")
fun depositProducerFactory(@Qualifier("consumerConfig") producerConfigs: Map<String, Any>): ProducerFactory<String, DepositEvent> {
return DefaultKafkaProducerFactory(producerConfigs)
}

@Bean("depositKafkaTemplate")
fun depositKafkaTemplate(@Qualifier("depositProducerFactory") producerFactory: ProducerFactory<String, DepositEvent>): KafkaTemplate<String, DepositEvent> {
return KafkaTemplate(producerFactory)
}

@Bean("depositKafkaListenerContainer")
@ConditionalOnBean(DepositKafkaListener::class)
fun configureDepositRequestEventListener(
fun depositListenerContainer(
listener: DepositKafkaListener,
@Qualifier("depositKafkaTemplate") template: KafkaTemplate<String, DepositEvent>,
@Qualifier("depositConsumerFactory") consumerFactory: ConsumerFactory<String, DepositEvent>
) {
): ConcurrentMessageListenerContainer<String, DepositEvent> {
val containerProps = ContainerProperties(Pattern.compile("deposit"))
containerProps.messageListener = listener
val container = ConcurrentMessageListenerContainer(consumerFactory, containerProps)
container.setBeanName("DepositKafkaListenerContainer")
container.commonErrorHandler = createConsumerErrorHandler(template, "deposit.DLT")
container.start()
return container
}

private fun createConsumerErrorHandler(kafkaTemplate: KafkaTemplate<*, *>, dltTopic: String): CommonErrorHandler {
Expand All @@ -224,5 +221,4 @@ class AccountantKafkaConfig {
}
return DefaultErrorHandler(recoverer, FixedBackOff(5_000, 20))
}

}
Original file line number Diff line number Diff line change
@@ -1,25 +1,52 @@
package co.nilin.opex.accountant.ports.walletproxy.config

import io.netty.channel.ChannelOption
import org.springframework.cloud.client.ServiceInstance
import org.springframework.cloud.client.loadbalancer.reactive.ReactiveLoadBalancer
import org.springframework.cloud.client.loadbalancer.reactive.ReactorLoadBalancerExchangeFilterFunction
import org.springframework.context.annotation.Bean
import org.springframework.context.annotation.Configuration
import org.springframework.http.client.reactive.ReactorClientHttpConnector
import org.springframework.web.reactive.function.client.WebClient
import org.zalando.logbook.Logbook
import org.zalando.logbook.netty.LogbookClientHandler
import reactor.netty.http.client.HttpClient
import reactor.netty.resources.ConnectionProvider
import java.time.Duration

@Configuration
class WebClientConfig {

@Bean
fun webClient(loadBalancerFactory: ReactiveLoadBalancer.Factory<ServiceInstance>, logbook: Logbook): WebClient {
val client = HttpClient.create().doOnConnected { it.addHandlerLast(LogbookClientHandler(logbook)) }
fun webClient(
loadBalancerFactory: ReactiveLoadBalancer.Factory<ServiceInstance>,
logbook: Logbook
): WebClient {

val connectionProvider = ConnectionProvider.builder("accountant-wallet")
.maxIdleTime(Duration.ofSeconds(20))
.maxLifeTime(Duration.ofMinutes(5))
.pendingAcquireTimeout(Duration.ofSeconds(5))
.evictInBackground(Duration.ofSeconds(30))
.lifo()
.build()

val client = HttpClient.create(connectionProvider)
.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 3000)
.responseTimeout(Duration.ofSeconds(10))
.keepAlive(true)
.doOnConnected {
it.addHandlerLast(LogbookClientHandler(logbook))
}

return WebClient.builder()
//.clientConnector(ReactorClientHttpConnector(client))
.filter(ReactorLoadBalancerExchangeFilterFunction(loadBalancerFactory, emptyList()))
.clientConnector(ReactorClientHttpConnector(client))
.filter(
ReactorLoadBalancerExchangeFilterFunction(
loadBalancerFactory,
emptyList()
)
)
.build()
}

}
}
Original file line number Diff line number Diff line change
Expand Up @@ -53,12 +53,17 @@ class RateLimitConfig(
return ReactiveSecurityContextHolder.getContext()
.mapNotNull { it.authentication }
.filter { it.isAuthenticated }
.flatMap { auth ->
applyRateLimit(auth.name, exchange, chain, groupId)
.map { auth ->
Mono.defer {
applyRateLimit(auth.name, exchange, chain, groupId)
}
}
.switchIfEmpty(
chain.filter(exchange)
.defaultIfEmpty(
Mono.defer {
chain.filter(exchange)
}
)
.flatMap { it }
}

private fun applyRateLimit(
Expand Down
Loading
Loading