Skip to content
Merged

Dev #711

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

}
26 changes: 25 additions & 1 deletion docker-compose.local.yml
Original file line number Diff line number Diff line change
Expand Up @@ -36,9 +36,33 @@ services:
postgres-otp:
ports:
- "127.0.0.1:5462:5432"
postgres-accountant:
ports:
- "5432:5432"
postgres-eventlog:
ports:
- "5433:5432"
postgres-auth:
ports:
- "5434:5432"
postgres-wallet:
ports:
- "5435:5432"
postgres-api:
ports:
- "5436:5432"
postgres-market:
ports:
- "127.0.0.1:5438:5432"
- "5438:5432"
postgres-bc-gateway:
ports:
- "5437:5432"
postgres-matching-gateway:
ports:
- "5439:5432"
postgres-profile:
ports:
- "5440:5432"
accountant:
ports:
- "127.0.0.1:8089:8080"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,13 +2,15 @@ package co.nilin.opex.market.core.event

import co.nilin.opex.market.core.inout.OrderStatus
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 @@ -7,7 +7,6 @@ import co.nilin.opex.market.ports.kafka.listener.consumer.TradeKafkaListener
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 @@ -52,34 +51,34 @@ class KafkaConsumerConfig {
return DefaultKafkaConsumerFactory(consumerConfigs)
}

@Autowired
@Bean("marketTradeKafkaListenerContainer")
@ConditionalOnBean(TradeKafkaListener::class)
fun configureTradeListener(
fun tradeListenerContainer(
tradeListener: TradeKafkaListener,
template: KafkaTemplate<String?, RichTrade>,
@Qualifier("richTradeConsumerFactory") consumerFactory: ConsumerFactory<String, RichTrade>
) {
): ConcurrentMessageListenerContainer<String, RichTrade> {
val containerProps = ContainerProperties(Pattern.compile("richTrade"))
containerProps.messageListener = tradeListener
val container = ConcurrentMessageListenerContainer(consumerFactory, containerProps)
container.setBeanName("marketTradeKafkaListenerContainer")
container.commonErrorHandler = createConsumerErrorHandler(template, "richTrade.DLT")
container.start()
return container
}

@Autowired
@Bean("marketOrderKafkaListenerContainer")
@ConditionalOnBean(OrderKafkaListener::class)
fun configureOrderListener(
fun orderListenerContainer(
orderListener: OrderKafkaListener,
template: KafkaTemplate<String?, RichOrderEvent>,
@Qualifier("richOrderConsumerFactory") consumerFactory: ConsumerFactory<String, RichOrderEvent>
) {
): ConcurrentMessageListenerContainer<String, RichOrderEvent> {
val containerProps = ContainerProperties(Pattern.compile("richOrder"))
containerProps.messageListener = orderListener
val container = ConcurrentMessageListenerContainer(consumerFactory, containerProps)
container.setBeanName("marketOrderKafkaListenerContainer")
container.commonErrorHandler = createConsumerErrorHandler(template, "richOrder.DLT")
container.start()
return container
}

private fun createConsumerErrorHandler(kafkaTemplate: KafkaTemplate<*, *>, dltTopic: String): CommonErrorHandler {
Expand Down
Loading
Loading