diff --git a/common/src/main/kotlin/co/nilin/opex/common/OpexError.kt b/common/src/main/kotlin/co/nilin/opex/common/OpexError.kt index 6edb3ba9d..109a6103d 100644 --- a/common/src/main/kotlin/co/nilin/opex/common/OpexError.kt +++ b/common/src/main/kotlin/co/nilin/opex/common/OpexError.kt @@ -32,6 +32,8 @@ enum class OpexError(val code: Int, val message: String?, val status: HttpStatus WithdrawAmountExceeds(2006, "The requested withdraw amount exceeds your daily limit", HttpStatus.BAD_REQUEST), WithdrawLimitConfigNotFound(2007, "Withdraw limit config not found", HttpStatus.NOT_FOUND), // code 3000: matching-engine + TemporaryInUnavailable(3000, "Transaction is temporarily unavailable. Please try again later.", HttpStatus.SERVICE_UNAVAILABLE), + // code 4000: matching-gateway SubmitOrderForbiddenByAccountant(4001, null, HttpStatus.BAD_REQUEST), diff --git a/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/bl/ExchangeEventHandler.kt b/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/bl/ExchangeEventHandler.kt index e967cdebc..7df18abaf 100644 --- a/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/bl/ExchangeEventHandler.kt +++ b/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/bl/ExchangeEventHandler.kt @@ -27,12 +27,9 @@ class ExchangeEventHandler( EventDispatcher.register(OrderBookPublishedEvent::class.java, localHandler) } - val handler: (CoreEvent) -> Unit = { - CoroutineScope(AppSchedulers.generalExecutor).launch { - eventsSubmitter.submit(it) - } + val handler: suspend (CoreEvent) -> Unit = { + eventsSubmitter.submit(it) } - val localHandler: (OrderBookPublishedEvent) -> Unit = { CoroutineScope(AppSchedulers.generalExecutor).launch { orderBookPersister.storeLastState(it.persistentOrderBook) diff --git a/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/bl/OrderBooks.kt b/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/bl/OrderBooks.kt index 95f0fc30a..99c2a0e3d 100644 --- a/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/bl/OrderBooks.kt +++ b/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/bl/OrderBooks.kt @@ -1,29 +1,87 @@ package co.nilin.opex.matching.engine.app.bl +import co.nilin.opex.common.OpexError +import co.nilin.opex.matching.engine.core.engine.SimpleOrderBook import co.nilin.opex.matching.engine.core.factory.OrderBookFactory -import co.nilin.opex.matching.engine.core.model.OrderBook import co.nilin.opex.matching.engine.core.model.Pair import co.nilin.opex.matching.engine.core.model.PersistentOrderBook +import co.nilin.opex.matching.engine.core.spi.OrderBookStore +import java.util.concurrent.ConcurrentHashMap -object OrderBooks { +object OrderBooks : OrderBookStore { - private val orderBooks = mutableMapOf() + private val orderBooks = + ConcurrentHashMap() fun createOrderBook(pair: String) { - println("Going to add order book:" + pair + ", current order books#" + orderBooks.size) - if (orderBooks.containsKey(pair)) - throw IllegalArgumentException("$pair has an order book right now!") + println( + "Going to add order book: $pair, " + + "current order books#${orderBooks.size}" + ) + val pairs = pair.split("_") - orderBooks[pair] = OrderBookFactory.createOrderBook(Pair(pairs[0], pairs[1])) - println("order book:" + pair + " added, current order books#" + orderBooks.size) + + require(pairs.size == 2) { + "Invalid pair format: $pair" + } + + val newBook = OrderBookFactory.createOrderBook( + Pair(pairs[0], pairs[1]) + ) + + val previous = orderBooks.putIfAbsent( + pair, + newBook + ) + + require(previous == null) { + "$pair has an order book right now!" + } + + println( + "Order book: $pair added, " + + "current order books#${orderBooks.size}" + ) } - fun reloadOrderBook(orderBook: PersistentOrderBook) { - orderBooks["${orderBook.pair.leftSideName}_${orderBook.pair.rightSideName}"] = - OrderBookFactory.createOrderBook(orderBook) + fun reloadOrderBook( + persistentOrderBook: PersistentOrderBook + ) { + val pairKey = + "${persistentOrderBook.pair.leftSideName}_" + + persistentOrderBook.pair.rightSideName + + orderBooks[pairKey] = + OrderBookFactory.createOrderBook( + persistentOrderBook + ) } - fun lookupOrderBook(pair: String): OrderBook { - return orderBooks[pair] ?: throw IllegalArgumentException("No orderbook for $pair") + override fun lookupOrderBook( + pair: String + ): SimpleOrderBook { + return orderBooks[pair] + ?: throw IllegalArgumentException( + "No order book for $pair" + ) } + + override fun replace( + pairKey: String, + expected: SimpleOrderBook, + replacement: SimpleOrderBook + ): Boolean { + synchronized(orderBooks) { + val current = orderBooks[pairKey] + + if (current !== expected) { + throw OpexError.InternalServerError.exception() + } + + orderBooks[pairKey] = replacement + return true + } + + } + } \ No newline at end of file diff --git a/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/config/AppConfig.kt b/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/config/AppConfig.kt index 246c936ee..cc7975b98 100644 --- a/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/config/AppConfig.kt +++ b/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/config/AppConfig.kt @@ -4,6 +4,8 @@ import co.nilin.opex.matching.engine.app.bl.ExchangeEventHandler import co.nilin.opex.matching.engine.app.bl.OrderBooks import co.nilin.opex.matching.engine.app.listener.MatchingEngineEventListener import co.nilin.opex.matching.engine.app.listener.OrderListener +import co.nilin.opex.matching.engine.core.engine.MatchingEngineRecoveryManager +import co.nilin.opex.matching.engine.core.engine.OrderCommandProcessor import co.nilin.opex.matching.engine.core.model.PersistentOrderBook import co.nilin.opex.matching.engine.core.spi.OrderBookPersister import co.nilin.opex.matching.engine.ports.kafka.listener.consumer.EventKafkaListener @@ -17,7 +19,7 @@ import org.springframework.context.annotation.Bean import org.springframework.context.annotation.Configuration @Configuration -class AppConfig { +class AppConfig(private val recoveryManager: MatchingEngineRecoveryManager, private val orderCommandProcessor: OrderCommandProcessor) { @Autowired private lateinit var symbols: List @@ -55,7 +57,7 @@ class AppConfig { @Bean fun orderListener(): OrderListener { - return OrderListener() + return OrderListener(orderCommandProcessor,recoveryManager) } @Autowired diff --git a/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/config/InitializeService.kt b/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/config/InitializeService.kt deleted file mode 100644 index 7caf5f262..000000000 --- a/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/config/InitializeService.kt +++ /dev/null @@ -1,14 +0,0 @@ -package co.nilin.opex.matching.engine.app.config - -import org.springframework.beans.factory.annotation.Value -import org.springframework.context.annotation.Bean -import org.springframework.context.annotation.Configuration - -@Configuration -class InitializeService { - - @Bean("symbols") - fun getSymbols(@Value("\${app.symbols}") symbols: String): List { - return symbols.split(",").map { it.trim() }.map { it.uppercase() } - } -} diff --git a/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/config/MatchingEngineCoreConfiguration.kt b/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/config/MatchingEngineCoreConfiguration.kt new file mode 100644 index 000000000..44c3ef0e1 --- /dev/null +++ b/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/config/MatchingEngineCoreConfiguration.kt @@ -0,0 +1,40 @@ +package co.nilin.opex.matching.engine.app.config + + +import co.nilin.opex.matching.engine.app.bl.OrderBooks +import co.nilin.opex.matching.engine.core.engine.MatchingEngineRecoveryManager +import co.nilin.opex.matching.engine.core.engine.OrderBookTransitionPreparer +import co.nilin.opex.matching.engine.core.engine.OrderCommandProcessor +import co.nilin.opex.matching.engine.core.spi.OrderBookStore +import co.nilin.opex.matching.engine.core.spi.OrderBookTransitionPublisher +import org.springframework.context.annotation.Bean +import org.springframework.context.annotation.Configuration + +@Configuration +class MatchingEngineCoreConfiguration { + + @Bean + fun matchingEngineRecoveryManager(): MatchingEngineRecoveryManager = + MatchingEngineRecoveryManager() + + @Bean + fun orderBookTransitionPreparer() = + OrderBookTransitionPreparer() + + @Bean + fun orderBookStore(): OrderBookStore = + OrderBooks + + @Bean + fun orderCommandProcessor( + transitionPreparer: OrderBookTransitionPreparer, + transitionPublisher: OrderBookTransitionPublisher, + recoveryManager: MatchingEngineRecoveryManager, + orderBookStore: OrderBookStore + ) = OrderCommandProcessor( + transitionPreparer = transitionPreparer, + transitionPublisher = transitionPublisher, + recoveryManager = recoveryManager, + orderBookStore = orderBookStore + ) +} \ No newline at end of file diff --git a/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/config/MatchingEngineInitializer.kt b/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/config/MatchingEngineInitializer.kt new file mode 100644 index 000000000..1e03647d1 --- /dev/null +++ b/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/config/MatchingEngineInitializer.kt @@ -0,0 +1,23 @@ +package co.nilin.opex.matching.engine.app.config + +import co.nilin.opex.matching.engine.core.engine.MatchingEngineRecoveryManager +import org.springframework.beans.factory.annotation.Value +import org.springframework.boot.context.event.ApplicationReadyEvent +import org.springframework.context.annotation.Bean +import org.springframework.context.annotation.Configuration +import org.springframework.context.event.EventListener + +@Configuration +class MatchingEngineInitializer( + private val recoveryManager: MatchingEngineRecoveryManager +) { + @EventListener(ApplicationReadyEvent::class) + fun initialize() { + recoveryManager.markRunning() + } + + @Bean("symbols") + fun getSymbols(@Value("\${app.symbols}") symbols: String): List { + return symbols.split(",").map { it.trim() }.map { it.uppercase() } + } +} \ No newline at end of file diff --git a/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/listener/OrderListener.kt b/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/listener/OrderListener.kt index b19ab8291..b4174ed8d 100644 --- a/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/listener/OrderListener.kt +++ b/matching-engine/matching-engine-app/src/main/kotlin/co/nilin/opex/matching/engine/app/listener/OrderListener.kt @@ -1,11 +1,19 @@ package co.nilin.opex.matching.engine.app.listener import co.nilin.opex.matching.engine.app.bl.OrderBooks -import co.nilin.opex.matching.engine.core.inout.* +import co.nilin.opex.matching.engine.core.engine.MatchingEngineRecoveryManager +import co.nilin.opex.matching.engine.core.engine.OrderCommandProcessor +import co.nilin.opex.matching.engine.core.inout.InputKafkaMetadata +import co.nilin.opex.matching.engine.core.inout.OrderRequestEvent +import co.nilin.opex.matching.engine.core.model.Pair +import co.nilin.opex.matching.engine.core.util.toCommand import co.nilin.opex.matching.engine.ports.kafka.listener.spi.OrderRequestEventListener import org.slf4j.LoggerFactory -class OrderListener : OrderRequestEventListener { +class OrderListener( + private val orderCommandProcessor: OrderCommandProcessor, + private val recoveryManager: MatchingEngineRecoveryManager +) : OrderRequestEventListener { private val logger = LoggerFactory.getLogger(OrderListener::class.java) @@ -13,37 +21,33 @@ class OrderListener : OrderRequestEventListener { return "OrderListener" } - override suspend fun onOrder(order: OrderRequestEvent, partition: Int, offset: Long, timestamp: Long) { - logger.info("OrderRequestEvent received. ${order::class.java.simpleName} ouid=${order.ouid}") - val orderBook = OrderBooks.lookupOrderBook( - order.pair.leftSideName + "_" - + order.pair.rightSideName + override suspend fun onOrder( + order: OrderRequestEvent, + partition: Int, + offset: Long, + timestamp: Long, + topic: String, + consumerGroupId: String + ) { + + logger.info( + "OrderRequestEvent received: type={}, ouid={}, partition={}, offset={}", + order::class.java.simpleName, + order.ouid, + partition, + offset ) - when (order) { - is OrderSubmitRequestEvent -> orderBook.handleNewOrderCommand( - OrderCreateCommand( - order.ouid, - order.uuid, - order.pair, - order.price, - order.quantity, - order.direction, - order.matchConstraint, - order.orderType - ) + val command = order.toCommand() + val currentBook = OrderBooks.lookupOrderBook(pairKey(order.pair)) + orderCommandProcessor.process( + command = command, + currentBook, + InputKafkaMetadata(topic, partition, offset, consumerGroupId) ) + } - is OrderCancelRequestEvent -> orderBook.handleCancelCommand( - OrderCancelCommand( - order.ouid, - order.uuid, - order.orderId, - order.pair - ) - ) - else -> logger.warn("Unknown event type of OrderRequestEvent") - } - } + private fun pairKey(pair: Pair): String = + "${pair.leftSideName}_${pair.rightSideName}" } \ No newline at end of file diff --git a/matching-engine/matching-engine-app/src/main/resources/application.yml b/matching-engine/matching-engine-app/src/main/resources/application.yml index 2e1778e5c..8d879b269 100644 --- a/matching-engine/matching-engine-app/src/main/resources/application.yml +++ b/matching-engine/matching-engine-app/src/main/resources/application.yml @@ -8,6 +8,9 @@ spring: bootstrap-servers: ${KAFKA_IP_PORT:localhost:9092} consumer: group-id: engine + producer: + properties: + max.request.size: 2097152 redis: host: ${REDIS_HOST:localhost} port: 6379 diff --git a/matching-engine/matching-engine-app/src/test/kotlin/co/nilin/opex/matching/engine/app/OrderBookEventEmitsUnitTest.kt b/matching-engine/matching-engine-app/src/test/kotlin/co/nilin/opex/matching/engine/app/OrderBookEventEmitsUnitTest.kt deleted file mode 100644 index dad59e359..000000000 --- a/matching-engine/matching-engine-app/src/test/kotlin/co/nilin/opex/matching/engine/app/OrderBookEventEmitsUnitTest.kt +++ /dev/null @@ -1,148 +0,0 @@ -package co.nilin.opex.matching.engine.app - -import co.nilin.opex.matching.engine.core.engine.SimpleOrderBook -import co.nilin.opex.matching.engine.core.eventh.EventDispatcher -import co.nilin.opex.matching.engine.core.eventh.events.OrderBookPublishedEvent -import co.nilin.opex.matching.engine.core.inout.OrderCancelCommand -import co.nilin.opex.matching.engine.core.inout.OrderCreateCommand -import co.nilin.opex.matching.engine.core.inout.OrderEditCommand -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 co.nilin.opex.matching.engine.core.model.PersistentOrderBook -import org.junit.jupiter.api.Assertions -import org.junit.jupiter.api.BeforeEach -import org.junit.jupiter.api.Test -import java.util.* - -class OrderBookEventEmitsUnitTest { - private val pair = co.nilin.opex.matching.engine.core.model.Pair("BTC", "USDT") - private val uuid = UUID.randomUUID().toString() - - private var persistentOrderBook: PersistentOrderBook? = null - - @BeforeEach - fun setup() { - val localHandler: (OrderBookPublishedEvent) -> Unit = { - persistentOrderBook = it.persistentOrderBook - } - EventDispatcher.register(OrderBookPublishedEvent::class.java, localHandler) - } - - @Test - fun givenOrderBook_whenOrderCreated_thenOrderBookEventPublished() { - //given - val orderBook = SimpleOrderBook(pair, false) - //when - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 1, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - //then - Assertions.assertNotNull(persistentOrderBook) - } - - - @Test - fun givenOrderBook_whenCancelOrder_thenOrderBookEventPublished() { - //given - val orderBook = SimpleOrderBook(pair, false) - val firstOrderId = UUID.randomUUID().toString() - val secondOrderId = UUID.randomUUID().toString() - - val firstOrder = orderBook.handleNewOrderCommand( - OrderCreateCommand( - firstOrderId, - uuid, - pair, - 2, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - secondOrderId, - uuid, - pair, - 1, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - persistentOrderBook = null - //when - orderBook.handleCancelCommand(OrderCancelCommand(firstOrderId, uuid, firstOrder!!.id()!!, pair)) - //then - Assertions.assertNotNull(persistentOrderBook) - } - - - @Test - fun givenOrderBook_whenEditOrder_thenOrderBookEventPublished() { - //given - val orderBook = SimpleOrderBook(pair, false) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 2, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - val secondOrder = orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 2, - 3, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 1, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - persistentOrderBook = null - //when - orderBook.handleEditCommand( - OrderEditCommand( - UUID.randomUUID().toString(), - uuid, - secondOrder!!.id()!!, - pair, - 3, - 2 - ) - ) - //then - Assertions.assertNotNull(persistentOrderBook) - } -} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/pom.xml b/matching-engine/matching-engine-core/pom.xml index e9c95993c..b784e614b 100644 --- a/matching-engine/matching-engine-core/pom.xml +++ b/matching-engine/matching-engine-core/pom.xml @@ -40,5 +40,18 @@ org.jetbrains.kotlinx kotlinx-coroutines-core + + co.nilin.opex.utility + error-handler + + + io.mockk + mockk + + + junit + junit + test + diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/engine/MatchingEngineRecoveryManager.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/engine/MatchingEngineRecoveryManager.kt new file mode 100644 index 000000000..8cf913da1 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/engine/MatchingEngineRecoveryManager.kt @@ -0,0 +1,65 @@ +package co.nilin.opex.matching.engine.core.engine + + +import co.nilin.opex.common.OpexError +import co.nilin.opex.matching.engine.core.model.MatchingEngineState +import co.nilin.opex.matching.engine.core.model.Pair +import org.slf4j.LoggerFactory +import java.util.concurrent.atomic.AtomicReference + +class MatchingEngineRecoveryManager { + + private val logger = + LoggerFactory.getLogger(MatchingEngineRecoveryManager::class.java) + + private val state = + AtomicReference(MatchingEngineState.STARTING) + + fun ensureRunning() { + val currentState = state.get() + + if (currentState != MatchingEngineState.RUNNING) { + throw OpexError.TemporaryInUnavailable.exception() + } + } + + fun markRunning() { + state.set(MatchingEngineState.RUNNING) + + logger.info("Matching engine state changed to RUNNING") + } + + fun enterRecovery( + cause: Throwable, + pair: Pair, + expectedSequence: Long + ) { + val previousState = + state.getAndSet(MatchingEngineState.RECOVERING) + + if (previousState == MatchingEngineState.RECOVERING) { + return + } + + logger.error( + "Matching engine entered recovery mode. " + + "pair={}_{}, expectedSequence={}", + pair.leftSideName, + pair.rightSideName, + expectedSequence, + cause + ) + } + + fun markFailed(cause: Throwable) { + state.set(MatchingEngineState.FAILED) + + logger.error( + "Matching engine entered FAILED state", + cause + ) + } + + fun currentState(): MatchingEngineState = + state.get() +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/engine/OrderBookTransitionPreparer.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/engine/OrderBookTransitionPreparer.kt new file mode 100644 index 000000000..2c4d99bc4 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/engine/OrderBookTransitionPreparer.kt @@ -0,0 +1,82 @@ +package co.nilin.opex.matching.engine.core.engine + +import co.nilin.opex.matching.engine.core.eventh.CollectingOrderBookEventSink +import co.nilin.opex.matching.engine.core.inout.OrderCancelCommand +import co.nilin.opex.matching.engine.core.inout.OrderCommand +import co.nilin.opex.matching.engine.core.inout.OrderCreateCommand +import co.nilin.opex.matching.engine.core.model.OrderBookDelta +import co.nilin.opex.matching.engine.core.model.PreparedCommandResult +import co.nilin.opex.matching.engine.core.model.PreparedStateTransition + +class OrderBookTransitionPreparer( +) { + + suspend fun prepare( + currentBook: SimpleOrderBook, + command: OrderCommand + ): PreparedCommandResult { + + val beforeSnapshot = + currentBook.snapshot() + + val eventCollector = + CollectingOrderBookEventSink() + + val workingBook = + SimpleOrderBook( + pair = currentBook.pair, + replayMode = false, + eventCollector + ) + + workingBook.rebuild( + beforeSnapshot + ) + + executeCommand( + book = workingBook, + command = command + ) + val nextSequence = + beforeSnapshot.sequence + 1 + + workingBook.sequence = + nextSequence + + + val afterSnapshot = + workingBook.snapshot() + + // We will implement this part next. +// TODO( +// "Compare before/after and create PreparedCommandResult" +// ) + + return PreparedCommandResult( + currentBook.pair, + commandId = command.ouid, + eventCollector.events(), + PreparedStateTransition( + beforeSnapshot.sequence, + workingBook.sequence, + afterSnapshot, + workingBook + ) + ) + } + + private suspend fun executeCommand( + book: SimpleOrderBook, + command: OrderCommand + ) { + when (command) { + is OrderCreateCommand -> + book.handleNewOrderCommand(command) + + is OrderCancelCommand -> + book.handleCancelCommand(command) + + else -> {} + } + } +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/engine/OrderCommandProcessor.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/engine/OrderCommandProcessor.kt new file mode 100644 index 000000000..aa5f2e0f4 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/engine/OrderCommandProcessor.kt @@ -0,0 +1,76 @@ +package co.nilin.opex.matching.engine.core.engine + +import co.nilin.opex.common.OpexError +import co.nilin.opex.matching.engine.core.inout.InputKafkaMetadata +import co.nilin.opex.matching.engine.core.inout.OrderCommand +import co.nilin.opex.matching.engine.core.model.* +import co.nilin.opex.matching.engine.core.spi.OrderBookStore +import co.nilin.opex.matching.engine.core.spi.OrderBookTransitionPublisher + +class OrderCommandProcessor( + private val transitionPreparer: OrderBookTransitionPreparer, + private val transitionPublisher: OrderBookTransitionPublisher, + private val recoveryManager: MatchingEngineRecoveryManager, + private val orderBookStore: OrderBookStore, + + ) { + + suspend fun process( + command: OrderCommand, + currentBook: SimpleOrderBook, + kafkaMetaData: InputKafkaMetadata + ) { + recoveryManager.ensureRunning() + + val pairKey = pairKey(command.pair) + + /* + * This executes matchInstantly(), putGtcInQueue(), + * cancellation and matching only on a working copy. + */ + val prepared = + transitionPreparer.prepare( + currentBook = currentBook, + command = command + ) + + /* + * One Kafka transaction: + * + * - all public events + * - order-book delta, if state changed + * - consumed input offset + 1 + */ + transitionPublisher.publish( + prepared = prepared, + inputMetadata = kafkaMetaData + ) + + /* + * Kafka has committed successfully. + * Now make the prepared state live. + */ + prepared.stateTransition?.let { transition -> + try { + orderBookStore.replace( + pairKey = pairKey, + expected = currentBook, + replacement = transition.preparedBook + ) + } catch (exception: Throwable) { + recoveryManager.enterRecovery( + cause = exception, + pair = command.pair, + expectedSequence = transition.nextSequence + ) + + throw OpexError.TemporaryInUnavailable.exception() + } + } + } + + private fun pairKey(pair: Pair): String = + "${pair.leftSideName}_${pair.rightSideName}" +} + + diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/engine/SimpleOrderBook.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/engine/SimpleOrderBook.kt index 5475734ee..91f2ff081 100644 --- a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/engine/SimpleOrderBook.kt +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/engine/SimpleOrderBook.kt @@ -1,18 +1,24 @@ package co.nilin.opex.matching.engine.core.engine -import co.nilin.opex.matching.engine.core.eventh.EventDispatcher +import co.nilin.opex.matching.engine.core.eventh.OrderBookEventSink import co.nilin.opex.matching.engine.core.eventh.events.* -import co.nilin.opex.matching.engine.core.inout.* +import co.nilin.opex.matching.engine.core.inout.OrderCancelCommand +import co.nilin.opex.matching.engine.core.inout.OrderCreateCommand +import co.nilin.opex.matching.engine.core.inout.RejectReason +import co.nilin.opex.matching.engine.core.inout.RequestedOperation import co.nilin.opex.matching.engine.core.model.* import exchange.core2.collections.art.LongAdaptiveRadixTreeMap import org.slf4j.LoggerFactory import java.util.* import java.util.concurrent.atomic.AtomicLong -class SimpleOrderBook(val pair: Pair, var replayMode: Boolean) : OrderBook { +class SimpleOrderBook( + val pair: Pair, + var replayMode: Boolean, + var eventSink: OrderBookEventSink +) : OrderBook { private val logger = LoggerFactory.getLogger(SimpleOrderBook::class.java) - val askOrders = LongAdaptiveRadixTreeMap() val bidOrders = LongAdaptiveRadixTreeMap() val orders = TreeMap() @@ -25,27 +31,28 @@ class SimpleOrderBook(val pair: Pair, var replayMode: Boolean) : OrderBook { var lastOrder: SimpleOrder? = null - override fun handleNewOrderCommand(orderCommand: OrderCreateCommand): Order? { + var sequence: Long = 0 + + override suspend fun handleNewOrderCommand(orderCommand: OrderCreateCommand): Order? { logNewOrder(orderCommand) val order = when (orderCommand.matchConstraint) { MatchConstraint.GTC -> { if (orderCommand.orderType == OrderType.MARKET_ORDER) { - if (!replayMode) { - EventDispatcher.emit( - RejectOrderEvent( - orderCommand.ouid, - orderCommand.uuid, - orderCommand.pair, - orderCommand.price, - orderCommand.quantity, - orderCommand.direction, - orderCommand.matchConstraint, - orderCommand.orderType, - RequestedOperation.PLACE_ORDER, - RejectReason.ORDER_TYPE_NOT_MATCHED_MATCHC - ) + emit( + RejectOrderEvent( + orderCommand.ouid, + orderCommand.uuid, + orderCommand.pair, + orderCommand.price, + orderCommand.quantity, + orderCommand.direction, + orderCommand.matchConstraint, + orderCommand.orderType, + RequestedOperation.PLACE_ORDER, + RejectReason.ORDER_TYPE_NOT_MATCHED_MATCHC ) - } + ) + return null } val order = SimpleOrder( @@ -62,22 +69,21 @@ class SimpleOrderBook(val pair: Pair, var replayMode: Boolean) : OrderBook { null, null ) - if (!replayMode) { - EventDispatcher.emit( - CreateOrderEvent( - orderCommand.ouid, - orderCommand.uuid, - order.id!!, - orderCommand.pair, - orderCommand.price, - orderCommand.quantity, - order.remainedQuantity(), - orderCommand.direction, - orderCommand.matchConstraint, - orderCommand.orderType - ) + emit( + CreateOrderEvent( + orderCommand.ouid, + orderCommand.uuid, + order.id!!, + orderCommand.pair, + orderCommand.price, + orderCommand.quantity, + order.remainedQuantity(), + orderCommand.direction, + orderCommand.matchConstraint, + orderCommand.orderType ) - } + ) + // try to match instantly val queueOrder = matchInstantly(order) // if remained quantity > 0 add to queue @@ -102,66 +108,59 @@ class SimpleOrderBook(val pair: Pair, var replayMode: Boolean) : OrderBook { null, null ) - if (!replayMode) { - EventDispatcher.emit( - CreateOrderEvent( - orderCommand.ouid, orderCommand.uuid, - order.id!!, orderCommand.pair, orderCommand.price, - orderCommand.quantity, order.remainedQuantity(), - orderCommand.direction, orderCommand.matchConstraint, orderCommand.orderType - ) + emit( + CreateOrderEvent( + orderCommand.ouid, orderCommand.uuid, + order.id!!, orderCommand.pair, orderCommand.price, + orderCommand.quantity, order.remainedQuantity(), + orderCommand.direction, orderCommand.matchConstraint, orderCommand.orderType ) - } + ) // try to match instantly val queueOrder = matchIocInstantly(order) - if (!replayMode) { - if (queueOrder.filledQuantity != queueOrder.quantity) { - EventDispatcher.emit( - CancelOrderEvent( - orderCommand.ouid, - orderCommand.uuid, - queueOrder.id!!, - orderCommand.pair, - order.price, - order.quantity, - order.remainedQuantity(), - order.direction, - order.matchConstraint, - order.orderType - ) + if (queueOrder.filledQuantity != queueOrder.quantity) { + emit( + CancelOrderEvent( + orderCommand.ouid, + orderCommand.uuid, + queueOrder.id!!, + orderCommand.pair, + order.price, + order.quantity, + order.remainedQuantity(), + order.direction, + order.matchConstraint, + order.orderType ) - } + ) } queueOrder } else -> { - if (!replayMode) { - EventDispatcher.emit( - RejectOrderEvent( - orderCommand.ouid, - orderCommand.uuid, - orderCommand.pair, - orderCommand.price, - orderCommand.quantity, - orderCommand.direction, - orderCommand.matchConstraint, - orderCommand.orderType, - RequestedOperation.PLACE_ORDER, - RejectReason.OPERATION_NOT_MATCHED_MATCHC - ) + emit( + RejectOrderEvent( + orderCommand.ouid, + orderCommand.uuid, + orderCommand.pair, + orderCommand.price, + orderCommand.quantity, + orderCommand.direction, + orderCommand.matchConstraint, + orderCommand.orderType, + RequestedOperation.PLACE_ORDER, + RejectReason.OPERATION_NOT_MATCHED_MATCHC ) - } + ) null } } lastOrder = order - EventDispatcher.emit(OrderBookPublishedEvent(persistent())) logCurrentState() return order } - override fun handleCancelCommand(orderCommand: OrderCancelCommand) { + override suspend fun handleCancelCommand(orderCommand: OrderCancelCommand) { logger.info( """ ---- CANCEL ${orderCommand.pair.leftSideName}-${orderCommand.pair.rightSideName} ---- @@ -174,18 +173,16 @@ class SimpleOrderBook(val pair: Pair, var replayMode: Boolean) : OrderBook { val simpleOrder = orders.entries.find { it.value.ouid == orderCommand.ouid } val order = simpleOrder?.value if (order == null /*check for userid*/) { - if (!replayMode) { - EventDispatcher.emit( - RejectOrderEvent( - orderCommand.ouid, - orderCommand.uuid, - orderCommand.orderId, - orderCommand.pair, - RequestedOperation.CANCEL_ORDER, - RejectReason.ORDER_NOT_FOUND - ) + emit( + RejectOrderEvent( + orderCommand.ouid, + orderCommand.uuid, + orderCommand.orderId, + orderCommand.pair, + RequestedOperation.CANCEL_ORDER, + RejectReason.ORDER_NOT_FOUND ) - } + ) return } else { orders.remove(simpleOrder.key) @@ -200,149 +197,17 @@ class SimpleOrderBook(val pair: Pair, var replayMode: Boolean) : OrderBook { bestAskOrder = newBestOrder } } - if (!replayMode) { - EventDispatcher.emit( - CancelOrderEvent( - orderCommand.ouid, orderCommand.uuid, - orderCommand.orderId, orderCommand.pair, - order.price, order.quantity, - order.remainedQuantity(), order.direction, - order.matchConstraint, order.orderType - ) + emit( + CancelOrderEvent( + orderCommand.ouid, orderCommand.uuid, + orderCommand.orderId, orderCommand.pair, + order.price, order.quantity, + order.remainedQuantity(), order.direction, + order.matchConstraint, order.orderType ) - } - EventDispatcher.emit(OrderBookPublishedEvent(persistent())) - logCurrentState() - } - - override fun handleEditCommand(orderCommand: OrderEditCommand): Order? { - val order = orders.remove(orderCommand.orderId) - if (order == null /*check for userid*/) { - if (!replayMode) { - EventDispatcher.emit( - RejectOrderEvent( - orderCommand.ouid, - orderCommand.uuid, - orderCommand.orderId, - orderCommand.pair, - RequestedOperation.EDIT_ORDER, - RejectReason.ORDER_NOT_FOUND - ) - ) - } - return order - } - if (order.direction == OrderDirection.BID) { - handleCancelOrder(order, bidOrders, bestBidOrder) { newBestOrder: SimpleOrder? -> - bestBidOrder = newBestOrder - } - } else { - handleCancelOrder(order, askOrders, bestAskOrder) { newBestOrder: SimpleOrder? -> - bestAskOrder = newBestOrder - } - } - val newOrder = SimpleOrder( - order.id, - orderCommand.ouid, - orderCommand.uuid, - orderCommand.price, - orderCommand.quantity, - order.matchConstraint, - order.orderType, - order.direction, - order.filledQuantity, - null, - null, - null ) - return when (order.matchConstraint) { - MatchConstraint.GTC -> { - if (!replayMode) { - EventDispatcher.emit( - UpdatedOrderEvent( - orderCommand.ouid, orderCommand.uuid, - order.id!!, orderCommand.pair, order.price, order.quantity, - orderCommand.price, orderCommand.quantity, order.remainedQuantity(), - order.direction, order.matchConstraint, order.orderType - ) - ) - } - // try to match instantly - val queueOrder = matchInstantly(newOrder) - //if remained quantity > 0 add to queue - if (queueOrder.filledQuantity != queueOrder.quantity) { - putGtcInQueue(queueOrder) - } - EventDispatcher.emit(OrderBookPublishedEvent(persistent())) - queueOrder - } - - MatchConstraint.IOC -> { - if (!replayMode) { - EventDispatcher.emit( - UpdatedOrderEvent( - orderCommand.ouid, - orderCommand.uuid, - order.id!!, - orderCommand.pair, - order.price, - order.quantity, - orderCommand.price, - orderCommand.quantity, - order.remainedQuantity(), - order.direction, - order.matchConstraint, - order.orderType - ) - ) - } - // try to match instantly - val queueOrder = matchIocInstantly(newOrder) - if (!replayMode) { - if (queueOrder.filledQuantity != queueOrder.quantity) { - EventDispatcher.emit( - CancelOrderEvent( - orderCommand.ouid, - orderCommand.uuid, - queueOrder.id!!, - orderCommand.pair, - order.price, - order.quantity, - order.remainedQuantity(), - order.direction, - order.matchConstraint, - order.orderType - ) - ) - } - } - EventDispatcher.emit(OrderBookPublishedEvent(persistent())) - queueOrder - } - - else -> { - if (!replayMode) { - EventDispatcher.emit( - RejectOrderEvent( - orderCommand.ouid, - orderCommand.uuid, - orderCommand.orderId, - orderCommand.pair, - orderCommand.price, - orderCommand.quantity, - order.direction, - order.matchConstraint, - order.orderType, - RequestedOperation.EDIT_ORDER, - RejectReason.OPERATION_NOT_MATCHED_MATCHC - ) - ) - } - null - } - } - + logCurrentState() } private fun handleCancelOrder( @@ -366,7 +231,7 @@ class SimpleOrderBook(val pair: Pair, var replayMode: Boolean) : OrderBook { setBestOrder(bestOrder.worse) } - private fun matchInstantly(order: SimpleOrder): SimpleOrder { + private suspend fun matchInstantly(order: SimpleOrder): SimpleOrder { if (order.direction == OrderDirection.BID) { return matchInstantly(order, bestAskOrder, askOrders, { makerPrice: Long -> makerPrice <= order.price @@ -382,7 +247,7 @@ class SimpleOrderBook(val pair: Pair, var replayMode: Boolean) : OrderBook { } } - private fun matchIocInstantly(order: SimpleOrder): SimpleOrder { + private suspend fun matchIocInstantly(order: SimpleOrder): SimpleOrder { if (order.direction == OrderDirection.BID) { return matchInstantly(order, bestAskOrder, askOrders, { makerPrice: Long -> order.orderType == OrderType.MARKET_ORDER || makerPrice <= order.price @@ -415,7 +280,7 @@ class SimpleOrderBook(val pair: Pair, var replayMode: Boolean) : OrderBook { } } - private fun matchInstantly( + private suspend fun matchInstantly( order: SimpleOrder, makerOrder: SimpleOrder?, queue: LongAdaptiveRadixTreeMap, @@ -433,28 +298,27 @@ class SimpleOrderBook(val pair: Pair, var replayMode: Boolean) : OrderBook { order.filledQuantity += instantMatchQuantity currentMaker.filledQuantity += instantMatchQuantity currentMaker.bucket!!.totalQuantity -= instantMatchQuantity - if (!replayMode) { - EventDispatcher.emit( - TradeEvent( - tradeCounter.incrementAndGet(), - pair, - order.ouid, - order.uuid, - order.id - ?: 0, - order.direction, - order.price, - order.remainedQuantity(), - currentMaker.ouid, - currentMaker.uuid, - currentMaker.id!!, - currentMaker.direction, - currentMaker.price, - currentMaker.remainedQuantity(), - instantMatchQuantity - ) + emit( + TradeEvent( + tradeCounter.incrementAndGet(), + pair, + order.ouid, + order.uuid, + order.id + ?: 0, + order.direction, + order.price, + order.remainedQuantity(), + currentMaker.ouid, + currentMaker.uuid, + currentMaker.id!!, + currentMaker.direction, + currentMaker.price, + currentMaker.remainedQuantity(), + instantMatchQuantity ) - } + ) + if (currentMaker.remainedQuantity() == 0L) { currentMaker.bucket!!.ordersCount-- } @@ -552,38 +416,42 @@ class SimpleOrderBook(val pair: Pair, var replayMode: Boolean) : OrderBook { return lastOrder } - private fun persistent(): PersistentOrderBook { - val persistent = PersistentOrderBook(pair) - persistent.lastOrder = lastOrder?.persistent() - persistent.orders = orders.values.map { order -> order.persistent() } - persistent.tradeCounter = tradeCounter.get() - return persistent - } fun rebuild(persistentOrderBook: PersistentOrderBook) { - persistentOrderBook.orders?.map { order -> - SimpleOrder( - order.id, - order.ouid, - order.uuid, - order.price, - order.quantity, - order.matchConstraint, - order.orderType, - order.direction, - order.filledQuantity, - null, - null, - null - ) - }?.filter { order -> - order.matchConstraint == MatchConstraint.GTC - }?.forEach { order -> putGtcInQueue(order) } + persistentOrderBook.orders + ?.filter { order -> + order.matchConstraint == MatchConstraint.GTC + } + ?.map { order -> + order.toSimpleOrderBook() + }?.forEach { order -> putGtcInQueue(order) } - orderCounter.set(persistentOrderBook.lastOrder?.id ?: 0) + orderCounter.set(persistentOrderBook.orderCounter) tradeCounter.set(persistentOrderBook.tradeCounter) + sequence = persistentOrderBook.sequence +// Reuse the existing in-memory order instance when available instead of creating a duplicate object. + lastOrder = persistentOrderBook.lastOrder?.let { lo -> + orders.getOrDefault(lo.id, lo.toSimpleOrderBook()) + } } + private fun PersistentOrder.toSimpleOrderBook(): SimpleOrder = + SimpleOrder( + id, + ouid, + uuid, + price, + quantity, + matchConstraint, + orderType, + direction, + filledQuantity, + null, + null, + null + ) + + private fun logNewOrder(orderCommand: OrderCreateCommand) { logger.info( """ @@ -613,4 +481,23 @@ class SimpleOrderBook(val pair: Pair, var replayMode: Boolean) : OrderBook { """.trimIndent() ) } + + private suspend fun emit(event: CoreEvent) { + if (!replayMode) { + eventSink.emit(event) + } + } + + fun snapshot(): PersistentOrderBook { + val snapshot = PersistentOrderBook(pair) + snapshot.sequence = sequence + snapshot.orderCounter = orderCounter.get() + snapshot.tradeCounter = tradeCounter.get() + snapshot.lastOrder = lastOrder?.persistent() + snapshot.orders = orders.values.map { order -> + order.persistent() + } + return snapshot + } + } \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/CollectingOrderBookEventSink.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/CollectingOrderBookEventSink.kt new file mode 100644 index 000000000..72f6f7002 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/CollectingOrderBookEventSink.kt @@ -0,0 +1,15 @@ +package co.nilin.opex.matching.engine.core.eventh + +import co.nilin.opex.matching.engine.core.eventh.events.CoreEvent + +class CollectingOrderBookEventSink : OrderBookEventSink { + + private val collectedEvents = mutableListOf() + + override suspend fun emit(event: CoreEvent) { + collectedEvents += event + } + + fun events(): List = + collectedEvents.toList() +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/EventDispatcher.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/EventDispatcher.kt index cc2d31cfc..c5e218a89 100644 --- a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/EventDispatcher.kt +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/EventDispatcher.kt @@ -13,16 +13,20 @@ object EventDispatcher { @JvmStatic fun register(type: Class, lambda: (T) -> Unit) = register(type, EventListener(lambda)) + @JvmStatic + fun register(type: Class, lambda: suspend (T) -> Unit) = register(type, EventListener(lambda)) + + @JvmStatic fun register(type: Class, listener: EventListener) { eventsHandler.getOrPut(type, { LinkedList() }).add(listener) } - fun emit(event: CoreEvent) { + suspend fun emit(event: CoreEvent) { var type: Class<*>? = event::class.java while (type != null) { eventsHandler[type]?.forEach { eventsHandler -> - kotlin.runCatching { + run { eventsHandler(event) } } @@ -30,10 +34,11 @@ object EventDispatcher { } } + open class EventListener( - val lambda: (T) -> Unit + val lambda: suspend (T) -> Unit ) { - operator fun invoke(event: Any) { + suspend operator fun invoke(event: Any) { lambda(event as T) } } diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/NoOpOrderBookEventSink.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/NoOpOrderBookEventSink.kt new file mode 100644 index 000000000..7f762d433 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/NoOpOrderBookEventSink.kt @@ -0,0 +1,10 @@ +package co.nilin.opex.matching.engine.core.eventh + +import co.nilin.opex.matching.engine.core.eventh.events.CoreEvent + +object NoOpOrderBookEventSink : OrderBookEventSink { + + override suspend fun emit(event: CoreEvent) { + // Live books do not publish events directly. + } +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/OrderBookEventSink.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/OrderBookEventSink.kt new file mode 100644 index 000000000..09c0913b5 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/OrderBookEventSink.kt @@ -0,0 +1,7 @@ +package co.nilin.opex.matching.engine.core.eventh + +import co.nilin.opex.matching.engine.core.eventh.events.CoreEvent + +fun interface OrderBookEventSink { + suspend fun emit(event: CoreEvent) +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/events/OrderBookDeltaEvent.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/events/OrderBookDeltaEvent.kt new file mode 100644 index 000000000..add0ff8da --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/eventh/events/OrderBookDeltaEvent.kt @@ -0,0 +1,44 @@ +package co.nilin.opex.matching.engine.core.eventh.events + +import co.nilin.opex.matching.engine.core.model.Pair +import co.nilin.opex.matching.engine.core.model.PersistentOrder + +/** + * Describes the minimal set of changes applied to the in-memory order book as a result of processing + * a single incoming command. It is intended to be published to a dedicated Kafka topic (compacted) + * and used for downstream projections and/or recovery replay. + */ +data class OrderBookDeltaEvent( + val version: Long, + val timestamp: Long, + val createdOrders: List, + val updatedOrders: List, + val removedOrderIds: List, + val bestAskId: Long?, + val bestBidId: Long?, + val note: String? = null, + // Keep CoreEvent contract (pair) +) : CoreEvent(Pair()) { + constructor( + pair: Pair, + version: Long, + timestamp: Long, + createdOrders: List, + updatedOrders: List, + removedOrderIds: List, + bestAskId: Long?, + bestBidId: Long?, + note: String? = null + ) : this( + version = version, + timestamp = timestamp, + createdOrders = createdOrders, + updatedOrders = updatedOrders, + removedOrderIds = removedOrderIds, + bestAskId = bestAskId, + bestBidId = bestBidId, + note = note + ) { + this.pair = pair + } +} diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/factory/OrderBookFactory.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/factory/OrderBookFactory.kt index e3a1010d3..d892108e2 100644 --- a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/factory/OrderBookFactory.kt +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/factory/OrderBookFactory.kt @@ -1,17 +1,61 @@ package co.nilin.opex.matching.engine.core.factory -import co.nilin.opex.matching.engine.core.model.OrderBook -import co.nilin.opex.matching.engine.core.model.PersistentOrderBook +import co.nilin.opex.matching.engine.core.engine.SimpleOrderBook +import co.nilin.opex.matching.engine.core.eventh.NoOpOrderBookEventSink +import co.nilin.opex.matching.engine.core.eventh.OrderBookEventSink +import co.nilin.opex.matching.engine.core.model.* object OrderBookFactory { - fun createOrderBook(pair: co.nilin.opex.matching.engine.core.model.Pair): OrderBook { - return co.nilin.opex.matching.engine.core.engine.SimpleOrderBook(pair, false) + + /** + * Creates an empty committed/live order book. + */ + fun createOrderBook(pair: Pair): SimpleOrderBook { + return SimpleOrderBook( + pair = pair, + replayMode = false, + eventSink = NoOpOrderBookEventSink + ) } - fun createOrderBook(persistentOrderBook: PersistentOrderBook): OrderBook { - val orderBook = co.nilin.opex.matching.engine.core.engine.SimpleOrderBook(persistentOrderBook.pair, true) + /** + * Restores a committed/live order book from a snapshot. + */ + fun createOrderBook( + persistentOrderBook: PersistentOrderBook + ): SimpleOrderBook { + val orderBook = SimpleOrderBook( + pair = persistentOrderBook.pair, + replayMode = true, + eventSink = NoOpOrderBookEventSink + ) + orderBook.rebuild(persistentOrderBook) orderBook.stopReplayMode() + return orderBook } + + /** + * Creates a temporary working book. + * + * CREATE and CANCEL commands are executed against this book. + * Generated events are sent to the provided collecting sink. + */ + fun createWorkingOrderBook( + snapshot: PersistentOrderBook, + eventSink: OrderBookEventSink + ): SimpleOrderBook { + + val workingBook = SimpleOrderBook( + pair = snapshot.pair, + replayMode = true, + eventSink = eventSink + ) + + workingBook.rebuild(snapshot) + workingBook.stopReplayMode() + + return workingBook + } } \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/InputKafkaMetaData.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/InputKafkaMetaData.kt new file mode 100644 index 000000000..e346af015 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/InputKafkaMetaData.kt @@ -0,0 +1,11 @@ +package co.nilin.opex.matching.engine.core.inout + +data class InputKafkaMetadata( + val topic: String, + val partition: Int, + val offset: Long, + val consumerGroupId: String +) { + val nextOffset: Long + get() = offset + 1 +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderCancelCommand.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderCancelCommand.kt index 87087d728..f074992d6 100644 --- a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderCancelCommand.kt +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderCancelCommand.kt @@ -2,4 +2,9 @@ package co.nilin.opex.matching.engine.core.inout import co.nilin.opex.matching.engine.core.model.Pair -class OrderCancelCommand(val ouid: String, val uuid: String, val orderId: Long, val pair: Pair) \ No newline at end of file +class OrderCancelCommand( + override val ouid: String, + override val uuid: String, + val orderId: Long, + override val pair: Pair +) : OrderCommand \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderCommand.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderCommand.kt new file mode 100644 index 000000000..5226b9b59 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderCommand.kt @@ -0,0 +1,8 @@ +package co.nilin.opex.matching.engine.core.inout +import co.nilin.opex.matching.engine.core.model.* + +sealed interface OrderCommand { + val ouid: String + val uuid: String + val pair: Pair +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderCreateCommand.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderCreateCommand.kt index da062871d..8af58d959 100644 --- a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderCreateCommand.kt +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderCreateCommand.kt @@ -7,12 +7,12 @@ import co.nilin.opex.matching.engine.core.model.Pair data class OrderCreateCommand( - val ouid: String, - val uuid: String, - val pair: Pair, + override val ouid: String, + override val uuid: String, + override val pair: Pair, val price: Long, val quantity: Long, val direction: OrderDirection, val matchConstraint: MatchConstraint, val orderType: OrderType -) \ No newline at end of file +): OrderCommand \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderEditCommand.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderEditCommand.kt deleted file mode 100644 index 1db450e3c..000000000 --- a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/inout/OrderEditCommand.kt +++ /dev/null @@ -1,12 +0,0 @@ -package co.nilin.opex.matching.engine.core.inout - -import co.nilin.opex.matching.engine.core.model.Pair - -data class OrderEditCommand( - val ouid: String, - val uuid: String, - val orderId: Long, - val pair: Pair, - val price: Long, - val quantity: Long -) \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/MatchingEngineState.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/MatchingEngineState.kt new file mode 100644 index 000000000..838147a28 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/MatchingEngineState.kt @@ -0,0 +1,8 @@ +package co.nilin.opex.matching.engine.core.model + +enum class MatchingEngineState { + STARTING, + RUNNING, + RECOVERING, + FAILED +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/OrderBook.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/OrderBook.kt index 07282860f..7f3643867 100644 --- a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/OrderBook.kt +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/OrderBook.kt @@ -2,14 +2,12 @@ package co.nilin.opex.matching.engine.core.model import co.nilin.opex.matching.engine.core.inout.OrderCancelCommand import co.nilin.opex.matching.engine.core.inout.OrderCreateCommand -import co.nilin.opex.matching.engine.core.inout.OrderEditCommand interface OrderBook { fun pair(): Pair fun startReplayMode() fun stopReplayMode() fun lastOrder(): Order? - fun handleNewOrderCommand(orderCommand: OrderCreateCommand): Order? - fun handleCancelCommand(orderCommand: OrderCancelCommand) - fun handleEditCommand(orderCommand: OrderEditCommand): Order? + suspend fun handleNewOrderCommand(orderCommand: OrderCreateCommand): Order? + suspend fun handleCancelCommand(orderCommand: OrderCancelCommand) } \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/OrderBookDelta.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/OrderBookDelta.kt new file mode 100644 index 000000000..5dd8a0a5d --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/OrderBookDelta.kt @@ -0,0 +1,17 @@ +package co.nilin.opex.matching.engine.core.model + + +data class OrderBookDelta( + val commandId: String, + val pair: Pair, + + val baseSequence: Long, + val nextSequence: Long, + + val upsertedOrders: List, + val removedOrderIds: Set, + + val orderCounter: Long, + val tradeCounter: Long, + val lastOrder: PersistentOrder? +) \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/PersistentOrderBook.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/PersistentOrderBook.kt index 49a7758be..d286b0e8c 100644 --- a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/PersistentOrderBook.kt +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/PersistentOrderBook.kt @@ -6,6 +6,8 @@ class PersistentOrderBook { var lastOrder: PersistentOrder? = null var orders: List? = emptyList() var tradeCounter: Long = 0 + var sequence: Long = 0 + var orderCounter: Long = 0 constructor() { } diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/PreparedCommandResult.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/PreparedCommandResult.kt new file mode 100644 index 000000000..5ac37e384 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/PreparedCommandResult.kt @@ -0,0 +1,10 @@ +package co.nilin.opex.matching.engine.core.model + +import co.nilin.opex.matching.engine.core.eventh.events.CoreEvent + +data class PreparedCommandResult( + val pair: Pair, + val commandId: String, + val events: List, + val stateTransition: PreparedStateTransition? +) \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/PreparedOrderBookTransition.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/PreparedOrderBookTransition.kt new file mode 100644 index 000000000..7f1cb4958 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/PreparedOrderBookTransition.kt @@ -0,0 +1,17 @@ +package co.nilin.opex.matching.engine.core.model + +import co.nilin.opex.matching.engine.core.engine.SimpleOrderBook +import co.nilin.opex.matching.engine.core.eventh.events.CoreEvent + +data class PreparedOrderBookTransition( + val commandId: String, + val pair: Pair, + + val baseSequence: Long, + val nextSequence: Long, + + val events: List, + val delta: OrderBookDelta, + + val preparedBook: SimpleOrderBook +) \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/PreparedStateTransition.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/PreparedStateTransition.kt new file mode 100644 index 000000000..d3378e434 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/model/PreparedStateTransition.kt @@ -0,0 +1,10 @@ +package co.nilin.opex.matching.engine.core.model + +import co.nilin.opex.matching.engine.core.engine.SimpleOrderBook + +data class PreparedStateTransition( + val baseSequence: Long, + val nextSequence: Long, + val snapshot: PersistentOrderBook, + val preparedBook: SimpleOrderBook +) \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/spi/OderBookStore.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/spi/OderBookStore.kt new file mode 100644 index 000000000..bec81bf3e --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/spi/OderBookStore.kt @@ -0,0 +1,15 @@ +package co.nilin.opex.matching.engine.core.spi + +import co.nilin.opex.matching.engine.core.engine.SimpleOrderBook + +interface OrderBookStore { + + fun lookupOrderBook(pairKey: String): SimpleOrderBook + + fun replace( + pairKey: String, + expected: SimpleOrderBook, + replacement: SimpleOrderBook + ): Boolean + +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/spi/OrderBookTransitionPublisher.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/spi/OrderBookTransitionPublisher.kt new file mode 100644 index 000000000..11a5251a2 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/spi/OrderBookTransitionPublisher.kt @@ -0,0 +1,12 @@ +package co.nilin.opex.matching.engine.core.spi + +import co.nilin.opex.matching.engine.core.inout.InputKafkaMetadata +import co.nilin.opex.matching.engine.core.model.PreparedCommandResult + +interface OrderBookTransitionPublisher { + + suspend fun publish( + prepared: PreparedCommandResult, + inputMetadata: InputKafkaMetadata + ) +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/util/convertor.kt b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/util/convertor.kt new file mode 100644 index 000000000..8b07207f3 --- /dev/null +++ b/matching-engine/matching-engine-core/src/main/kotlin/co/nilin/opex/matching/engine/core/util/convertor.kt @@ -0,0 +1,31 @@ +package co.nilin.opex.matching.engine.core.util + +import co.nilin.opex.matching.engine.core.inout.* + + fun OrderRequestEvent.toCommand(): OrderCommand = + when (this) { + is OrderSubmitRequestEvent -> + OrderCreateCommand( + ouid = ouid, + uuid = uuid, + pair = pair, + price = price, + quantity = quantity, + direction = direction, + matchConstraint = matchConstraint, + orderType = orderType + ) + + is OrderCancelRequestEvent -> + OrderCancelCommand( + ouid = ouid, + uuid = uuid, + orderId = orderId, + pair = pair + ) + + else -> + throw IllegalArgumentException( + "Unsupported order request type: ${this::class.java.name}" + ) + } \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/FakeOrderBookStore.kt b/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/FakeOrderBookStore.kt new file mode 100644 index 000000000..3fc001d07 --- /dev/null +++ b/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/FakeOrderBookStore.kt @@ -0,0 +1,44 @@ +package co.nilin.opex.matching.engine.core + +import co.nilin.opex.matching.engine.core.engine.SimpleOrderBook +import co.nilin.opex.matching.engine.core.spi.OrderBookStore + +class FakeOrderBookStore( + initialBook: SimpleOrderBook +) : OrderBookStore { + + var currentBook = initialBook + private set + + var replaceCount = 0 + + + var failure: RuntimeException? = null + + override fun lookupOrderBook( + pairKey: String + ): SimpleOrderBook = + currentBook + + override fun replace( + pairKey: String, + expected: SimpleOrderBook, + replacement: SimpleOrderBook + ): Boolean { + replaceCount++ + + if (currentBook !== expected) { + return false + } + + failure?.let { throw it } + + if (currentBook !== expected) { + throw IllegalStateException("Current book has changed") + } + + currentBook = replacement + return true + } + +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/FakeOrderBookTransitionPublisher.kt b/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/FakeOrderBookTransitionPublisher.kt new file mode 100644 index 000000000..b3ac25eed --- /dev/null +++ b/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/FakeOrderBookTransitionPublisher.kt @@ -0,0 +1,21 @@ +package co.nilin.opex.matching.engine.core + +import co.nilin.opex.matching.engine.core.inout.InputKafkaMetadata +import co.nilin.opex.matching.engine.core.model.PreparedCommandResult +import co.nilin.opex.matching.engine.core.spi.OrderBookTransitionPublisher + +class FakeOrderBookTransitionPublisher : OrderBookTransitionPublisher { + var published: PreparedCommandResult? = null + var publishCount = 0 + var failure: RuntimeException? = null + + override suspend fun publish( + prepared: PreparedCommandResult, + inputMetadata: InputKafkaMetadata + ) { + publishCount++ + failure?.let { throw it } + published = prepared + } + +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/engine/OrderBookTransitionPreparerTest.kt b/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/engine/OrderBookTransitionPreparerTest.kt new file mode 100644 index 000000000..1e68eca96 --- /dev/null +++ b/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/engine/OrderBookTransitionPreparerTest.kt @@ -0,0 +1,216 @@ +package co.nilin.opex.matching.engine.core.engine + +import co.nilin.opex.matching.engine.core.engine.OrderBookTransitionPreparer +import co.nilin.opex.matching.engine.core.engine.SimpleOrderBook +import co.nilin.opex.matching.engine.core.eventh.CollectingOrderBookEventSink +import co.nilin.opex.matching.engine.core.eventh.events.CreateOrderEvent +import co.nilin.opex.matching.engine.core.eventh.events.RejectOrderEvent +import co.nilin.opex.matching.engine.core.eventh.events.TradeEvent +import co.nilin.opex.matching.engine.core.inout.OrderCancelCommand +import co.nilin.opex.matching.engine.core.inout.OrderCreateCommand +import co.nilin.opex.matching.engine.core.model.* +import kotlinx.coroutines.runBlocking +import org.assertj.core.api.Assertions.assertThat +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertNotSame +import org.junit.jupiter.api.Test +import java.util.* + +class OrderBookTransitionPreparerTest { + + val pair = Pair("BTC", "USDT") + val userId = "user-1" + + private val transitionPreparer = + OrderBookTransitionPreparer() + + @Test + fun givenExistingOrder_whenTransitionPrepared_thenHistoricalEventIsNotEmittedAgain() { + runBlocking { + + val sourceBook = generateOrderBook() + + + sourceBook.handleNewOrderCommand( + generateBidCommand() + ) + + + val beforeSnapShot = sourceBook.snapshot() + val newIncomingCommand = generateBidCommand() + + val prepared = transitionPreparer.prepare(sourceBook, newIncomingCommand) + + assertEquals(beforeSnapShot.orderCounter, sourceBook.orderCounter.get()) + assertEquals(beforeSnapShot.tradeCounter, sourceBook.tradeCounter.get()) + assertEquals(beforeSnapShot.sequence, sourceBook.sequence) + assertNotSame(sourceBook, prepared.stateTransition?.preparedBook) + + + + assertThat(prepared.events.filterIsInstance().size).isEqualTo(1) + assertEquals(1, prepared.events.size) + assertThat((prepared.events.single() as CreateOrderEvent).ouid).isEqualTo(newIncomingCommand.ouid) + + + } + } + + @Test + fun givenExistingOrder_whenTransitionPrepared_thenMatchingHappenedInWorkingBookAndSourceBookRemainsUnchanged() { + runBlocking { + val sourceBook = generateOrderBook() + + + val existingCommand = generateAskCommand() + sourceBook.handleNewOrderCommand( + existingCommand + ) + + + val beforeSnapShot = sourceBook.snapshot() + val newIncomingCommand = generateBidCommand() + + + val prepared = transitionPreparer.prepare(sourceBook, newIncomingCommand) + + val workingBook = prepared.stateTransition?.preparedBook + assertEquals(beforeSnapShot.orderCounter, sourceBook.orderCounter.get()) + assertEquals(0, beforeSnapShot.tradeCounter) + + assertEquals(1, workingBook?.tradeCounter?.get()) + + assertEquals( + 2, workingBook?.orders?.filter { (_, order) -> order.ouid == newIncomingCommand.ouid }?.values?.sumOf( + SimpleOrder::filledQuantity + ) + ) + assertEquals( + 0, sourceBook.orders.filter { (_, order) -> order.ouid == existingCommand.ouid }.values.sumOf( + SimpleOrder::filledQuantity + ) + ) + + + val createEvents = + prepared.events.filterIsInstance() + + val tradeEvents = + prepared.events.filterIsInstance() + + assertThat(createEvents) + .hasSize(1) + + assertThat(createEvents.single().ouid) + .isEqualTo(newIncomingCommand.ouid) + + assertThat(tradeEvents) + .hasSize(1) + } + } + + @Test + fun givenExistingOrder_whenTransitionPreparedForCancellationCommand_thenSourceBookRemainsUnchanged() { + runBlocking { + val sourceBook = generateOrderBook() + + + val existingAskCommand1 = generateAskCommand() + val existingAskCommand2 = generateAskCommand() + + sourceBook.handleNewOrderCommand( + existingAskCommand1 + ) + sourceBook.handleNewOrderCommand( + existingAskCommand2 + ) + + + val beforeSnapShot = sourceBook.snapshot() + val cancelCommand = generateCancelCommand(existingAskCommand1.ouid) + + + val prepared = transitionPreparer.prepare(sourceBook, cancelCommand) + + val workingBook = prepared.stateTransition?.preparedBook + + assertEquals(2, sourceBook.orderCounter.get()) + assertEquals(2, beforeSnapShot.orderCounter) + + assertEquals(1, workingBook?.orders?.size) + + assertThat(sourceBook.orders.values.any { order -> order.ouid == cancelCommand.ouid }).isTrue() + assertThat(workingBook?.orders?.values?.any { order -> order.ouid == cancelCommand.ouid }).isFalse() + + + } + } + + @Test + fun givenExistingOrder_whenTransitionPreparedForInvalidCancellationCommand_thenEmitRejection() { + runBlocking { + val eventCollector = CollectingOrderBookEventSink() + val sourceBook = generateOrderBook(eventCollector) + + + val existingAskCommand = generateAskCommand() + + sourceBook.handleNewOrderCommand( + existingAskCommand + ) + + + val invalidCancelCommand = generateCancelCommand(UUID.randomUUID().toString()) + + + val prepared = transitionPreparer.prepare(sourceBook, invalidCancelCommand) + + val workingBookRejectionEvents = prepared.events.filterIsInstance() + val sourceBookRejectionEvents = eventCollector.events().filterIsInstance() + + + assertEquals(1, workingBookRejectionEvents.size) + assertEquals(0, sourceBookRejectionEvents.size) + + assertThat( + workingBookRejectionEvents.single().ouid + ).isEqualTo(invalidCancelCommand.ouid) + } + } + + + private fun generateBidCommand() = OrderCreateCommand( + ouid = UUID.randomUUID().toString(), + uuid = userId, + pair = pair, + price = 100, + quantity = 5, + direction = OrderDirection.BID, + matchConstraint = MatchConstraint.GTC, + orderType = OrderType.LIMIT_ORDER + ) + + private fun generateCancelCommand(ouid: String) = OrderCancelCommand( + ouid = ouid, + uuid = userId, + pair = pair, + orderId = Long.MAX_VALUE, + ) + + private fun generateAskCommand() = OrderCreateCommand( + ouid = UUID.randomUUID().toString(), + uuid = userId, + pair = pair, + price = 100, + quantity = 2, + direction = OrderDirection.ASK, + matchConstraint = MatchConstraint.GTC, + orderType = OrderType.LIMIT_ORDER + ) + + private fun generateOrderBook(eventCollector: CollectingOrderBookEventSink? = null) = SimpleOrderBook( + pair = pair, + replayMode = false, + eventSink = eventCollector ?: CollectingOrderBookEventSink() + ) +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/engine/OrderCommandProcessorTest.kt b/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/engine/OrderCommandProcessorTest.kt new file mode 100644 index 000000000..9bfcc6740 --- /dev/null +++ b/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/engine/OrderCommandProcessorTest.kt @@ -0,0 +1,228 @@ +package co.nilin.opex.matching.engine.core.engine + +import co.nilin.opex.common.OpexError +import co.nilin.opex.matching.engine.core.FakeOrderBookStore +import co.nilin.opex.matching.engine.core.FakeOrderBookTransitionPublisher +import co.nilin.opex.matching.engine.core.eventh.CollectingOrderBookEventSink +import co.nilin.opex.matching.engine.core.inout.InputKafkaMetadata +import co.nilin.opex.matching.engine.core.inout.OrderCancelCommand +import co.nilin.opex.matching.engine.core.inout.OrderCreateCommand +import co.nilin.opex.matching.engine.core.model.* +import io.mockk.coEvery +import io.mockk.mockk +import kotlinx.coroutines.runBlocking +import org.junit.jupiter.api.Assertions.* +import org.junit.jupiter.api.BeforeEach +import org.junit.jupiter.api.Test +import org.junit.jupiter.api.assertThrows +import java.util.* + +class OrderCommandProcessorTest { + private val pair = Pair("BTC", "USDT") + private val userId = "user-1" + private lateinit var preparer: OrderBookTransitionPreparer + private lateinit var publisher: FakeOrderBookTransitionPublisher + private lateinit var recoveryManager: MatchingEngineRecoveryManager + private lateinit var originalBook: SimpleOrderBook + private lateinit var orderBookStore: FakeOrderBookStore + private lateinit var processor: OrderCommandProcessor + + @BeforeEach + fun setup() { + preparer = OrderBookTransitionPreparer() + recoveryManager = MatchingEngineRecoveryManager() + originalBook = SimpleOrderBook(pair, false, CollectingOrderBookEventSink()) + orderBookStore = FakeOrderBookStore(originalBook) + publisher = FakeOrderBookTransitionPublisher() + processor = OrderCommandProcessor(preparer, publisher, recoveryManager, orderBookStore) + recoveryManager.markRunning() + + } + + @Test + fun givenValidCommand_whenPublishedSuccessfully_thenPreparedBookReplacesLiveBook() { + runBlocking { + + val command = generateBidCommand() + + val inputMetadata = generateMetaDta() + + processor.process( + command = command, + currentBook = originalBook, + kafkaMetaData = inputMetadata + ) + + assertEquals( + 1, + publisher.publishCount + ) + + assertEquals( + 1, + orderBookStore.replaceCount + ) + + assertNotSame( + originalBook, + orderBookStore.currentBook + ) + + assertSame( + publisher.published + ?.stateTransition + ?.preparedBook, + orderBookStore.currentBook + ) + } + + } + + + @Test + fun givenValidCommand_whenPublishedThrowsException_thenReplacementWillNeverCall() { + runBlocking { + val command = generateBidCommand() + + preparer = mockk() + processor = OrderCommandProcessor(preparer, publisher, recoveryManager, orderBookStore) + val inputMetadata = generateMetaDta() + + val prepared = PreparedCommandResult(pair, command.ouid, emptyList(), null) + + coEvery { + preparer.prepare( + originalBook, command + ) + }.returns(prepared) + + publisher.failure = RuntimeException("Kafka exception in publishment") + + recoveryManager.markRunning() + + + assertThrows(RuntimeException::class.java) { + runBlocking { + processor.process( + command = command, + currentBook = originalBook, + kafkaMetaData = inputMetadata + ) + } + } + + assertEquals( + 1, + publisher.publishCount + ) + + assertEquals( + 0, + orderBookStore.replaceCount + ) + } + } + + @Test + fun givenPublicationSucceeds_whenLiveBookReplacementThrows_thenProcessorEntersRecovery() { + val command = generateBidCommand() + val inputMetadata = generateMetaDta() + + val replacementFailure = + RuntimeException("Live-book replacement failed") + + orderBookStore.failure = + replacementFailure + + recoveryManager.markRunning() + + val thrown = assertThrows { + runBlocking { + processor.process( + command = command, + currentBook = originalBook, + kafkaMetaData = inputMetadata + ) + } + } + + // Publication completed before the local replacement. + assertEquals( + 1, + publisher.publishCount + ) + + // Replacement was attempted once. + assertEquals( + 1, + orderBookStore.replaceCount + ) + + // The local live book was not changed. + assertSame( + originalBook, + orderBookStore.currentBook + ) + + // Processor translated the internal replacement failure + // into the public temporary-unavailable error. + assertEquals( + OpexError.TemporaryInUnavailable.exception()::class, + thrown::class + ) + + // enterRecovery() means the engine must no longer be RUNNING. + assertThrows { + recoveryManager.ensureRunning() + } + + assertThrows { + recoveryManager.ensureRunning() + } + + assertEquals(MatchingEngineState.RECOVERING, recoveryManager.currentState()) + } + + + private fun generateBidCommand() = OrderCreateCommand( + ouid = UUID.randomUUID().toString(), + uuid = userId, + pair = pair, + price = 100, + quantity = 5, + direction = OrderDirection.BID, + matchConstraint = MatchConstraint.GTC, + orderType = OrderType.LIMIT_ORDER + ) + + private fun generateCancelCommand(ouid: String) = OrderCancelCommand( + ouid = ouid, + uuid = userId, + pair = pair, + orderId = Long.MAX_VALUE, + ) + + private fun generateAskCommand() = OrderCreateCommand( + ouid = UUID.randomUUID().toString(), + uuid = userId, + pair = pair, + price = 100, + quantity = 2, + direction = OrderDirection.ASK, + matchConstraint = MatchConstraint.GTC, + orderType = OrderType.LIMIT_ORDER + ) + + private fun generateOrderBook(eventCollector: CollectingOrderBookEventSink? = null) = SimpleOrderBook( + pair = pair, + replayMode = false, + eventSink = eventCollector ?: CollectingOrderBookEventSink() + ) + + private fun generateMetaDta() = InputKafkaMetadata( + topic = "order-requests", + partition = 0, + offset = 10, + consumerGroupId = "matching-engine" + ) +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/engine/SimpleOrderBookSnapshotTest.kt b/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/engine/SimpleOrderBookSnapshotTest.kt new file mode 100644 index 000000000..1e29a460e --- /dev/null +++ b/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/engine/SimpleOrderBookSnapshotTest.kt @@ -0,0 +1,112 @@ +package co.nilin.opex.matching.engine.core.engine + +import co.nilin.opex.matching.engine.core.engine.SimpleOrderBook +import co.nilin.opex.matching.engine.core.eventh.CollectingOrderBookEventSink +import co.nilin.opex.matching.engine.core.inout.OrderCreateCommand +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 kotlinx.coroutines.runBlocking +import co.nilin.opex.matching.engine.core.model.Pair +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertNotNull +import org.junit.jupiter.api.Assertions.assertNotSame +import org.junit.jupiter.api.Test + +class SimpleOrderBookSnapshotTest { + private val pair = Pair("BTC", "USDT") + private val userId = "user-1" + + @Test + fun givenOrderBook_whenSnapshotRebuilt_thenCanonicalStateIsRestored(): Unit = + runBlocking { + val sourceBook = SimpleOrderBook( + pair = pair, + replayMode = false, + eventSink = CollectingOrderBookEventSink() + ) + + sourceBook.handleNewOrderCommand( + OrderCreateCommand( + ouid = "bid-1", + uuid = userId, + pair = pair, + price = 100, + quantity = 5, + direction = OrderDirection.BID, + matchConstraint = MatchConstraint.GTC, + orderType = OrderType.LIMIT_ORDER + ) + ) + + sourceBook.handleNewOrderCommand( + OrderCreateCommand( + ouid = "ask-1", + uuid = userId, + pair = pair, + price = 110, + quantity = 3, + direction = OrderDirection.ASK, + matchConstraint = MatchConstraint.GTC, + orderType = OrderType.LIMIT_ORDER + ) + ) + + sourceBook.sequence = 12 + + val snapshot = sourceBook.snapshot() + + val restoredBook = SimpleOrderBook( + pair = pair, + replayMode = true, + eventSink = CollectingOrderBookEventSink() + ) + + restoredBook.rebuild(snapshot) + restoredBook.stopReplayMode() + + assertNotSame(sourceBook, restoredBook) + assertEquals(sourceBook.sequence, restoredBook.sequence) + assertEquals( + sourceBook.orderCounter.get(), + restoredBook.orderCounter.get() + ) + assertEquals( + sourceBook.tradeCounter.get(), + restoredBook.tradeCounter.get() + ) + + assertEquals( + sourceBook.orders.size, + restoredBook.orders.size + ) + + assertEquals( + sourceBook.bidOrders.entriesList().size, + restoredBook.bidOrders.entriesList().size + ) + + assertEquals( + sourceBook.askOrders.entriesList().size, + restoredBook.askOrders.entriesList().size + ) + + assertNotNull(restoredBook.bestBidOrder) + assertNotNull(restoredBook.bestAskOrder) + + assertEquals( + sourceBook.bestBidOrder!!.price, + restoredBook.bestBidOrder!!.price + ) + + assertEquals( + sourceBook.bestAskOrder!!.price, + restoredBook.bestAskOrder!!.price + ) + + assertEquals( + sourceBook.lastOrder!!.ouid, + restoredBook.lastOrder!!.ouid + ) + } +} \ No newline at end of file diff --git a/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/SimpleOrderBookUnitTest.kt b/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/engine/SimpleOrderBookUnitTest.kt similarity index 57% rename from matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/SimpleOrderBookUnitTest.kt rename to matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/engine/SimpleOrderBookUnitTest.kt index 8b67edf10..52cb412b7 100644 --- a/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/SimpleOrderBookUnitTest.kt +++ b/matching-engine/matching-engine-core/src/test/kotlin/co/nilin/opex/matching/engine/core/engine/SimpleOrderBookUnitTest.kt @@ -1,14 +1,15 @@ -package co.nilin.opex.matching.engine.core +package co.nilin.opex.matching.engine.core.engine import co.nilin.opex.matching.engine.core.engine.SimpleOrderBook +import co.nilin.opex.matching.engine.core.eventh.CollectingOrderBookEventSink import co.nilin.opex.matching.engine.core.inout.OrderCancelCommand import co.nilin.opex.matching.engine.core.inout.OrderCreateCommand -import co.nilin.opex.matching.engine.core.inout.OrderEditCommand 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 co.nilin.opex.matching.engine.core.model.SimpleOrder import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.runBlocking import org.junit.jupiter.api.Assertions import org.junit.jupiter.api.Test import java.util.* @@ -19,9 +20,9 @@ class SimpleOrderBookUnitTest { private val uuid = UUID.randomUUID().toString() @Test - fun givenEmptyOrderBook_whenGtcBidLimitOrderCreated_then1BucketWithSize1() { + fun givenEmptyOrderBook_whenGtcBidLimitOrderCreated_then1BucketWithSize1(): Unit = runBlocking { //given - val orderBook = SimpleOrderBook(pair, false) + val orderBook = createOrderBook() //when val order = orderBook.handleNewOrderCommand( OrderCreateCommand( @@ -42,9 +43,9 @@ class SimpleOrderBookUnitTest { } @Test - fun givenOrderBookWithBidOrders_whenGtcBidLimitOrderWithSamePriceCreated_then() { + fun givenOrderBookWithBidOrders_whenGtcBidLimitOrderWithSamePriceCreated_then(): Unit = runBlocking { //given - val orderBook = SimpleOrderBook(pair, false) + val orderBook = createOrderBook() orderBook.handleNewOrderCommand( OrderCreateCommand( UUID.randomUUID().toString(), @@ -83,91 +84,93 @@ class SimpleOrderBookUnitTest { } @Test - fun givenOrderBookWithBidOrders_whenGtcBidLimitOrderWithLowerPriceCreated_thenBestOrderNotChange() { - //given - val orderBook = SimpleOrderBook(pair, false) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 2, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - val bestBidOrder = orderBook.bestBidOrder - //when - val order: SimpleOrder = + fun givenOrderBookWithBidOrders_whenGtcBidLimitOrderWithLowerPriceCreated_thenBestOrderNotChange(): Unit = + runBlocking { + //given + val orderBook = createOrderBook() orderBook.handleNewOrderCommand( OrderCreateCommand( UUID.randomUUID().toString(), uuid, pair, - 1, + 2, 1, OrderDirection.BID, MatchConstraint.GTC, OrderType.LIMIT_ORDER ) - ) as SimpleOrder - //then - Assertions.assertEquals(orderBook.bidOrders.entriesList().size, 2) - Assertions.assertEquals(orderBook.bestBidOrder, bestBidOrder) - Assertions.assertEquals(bestBidOrder!!.worse, order) - Assertions.assertEquals(order.better, bestBidOrder) - Assertions.assertEquals(orderBook.bidOrders.get(order.price).lastOrder, order) - Assertions.assertEquals(orderBook.bidOrders.get(order.price).totalQuantity, 1) - Assertions.assertEquals(orderBook.bidOrders.get(order.price).ordersCount, 1) - } + ) + val bestBidOrder = orderBook.bestBidOrder + //when + val order: SimpleOrder = + orderBook.handleNewOrderCommand( + OrderCreateCommand( + UUID.randomUUID().toString(), + uuid, + pair, + 1, + 1, + OrderDirection.BID, + MatchConstraint.GTC, + OrderType.LIMIT_ORDER + ) + ) as SimpleOrder + //then + Assertions.assertEquals(orderBook.bidOrders.entriesList().size, 2) + Assertions.assertEquals(orderBook.bestBidOrder, bestBidOrder) + Assertions.assertEquals(bestBidOrder!!.worse, order) + Assertions.assertEquals(order.better, bestBidOrder) + Assertions.assertEquals(orderBook.bidOrders.get(order.price).lastOrder, order) + Assertions.assertEquals(orderBook.bidOrders.get(order.price).totalQuantity, 1) + Assertions.assertEquals(orderBook.bidOrders.get(order.price).ordersCount, 1) + } @Test - fun givenOrderBookWithBidOrders_whenGtcBidLimitOrderWithHigherPriceCreated_thenBestOrderChanged() { - //given - val orderBook = SimpleOrderBook(pair, false) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 1, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - val bestBidOrder = orderBook.bestBidOrder - //when - val order: SimpleOrder = + fun givenOrderBookWithBidOrders_whenGtcBidLimitOrderWithHigherPriceCreated_thenBestOrderChanged(): Unit = + runBlocking { + //given + val orderBook = createOrderBook() orderBook.handleNewOrderCommand( OrderCreateCommand( UUID.randomUUID().toString(), uuid, pair, - 2, + 1, 1, OrderDirection.BID, MatchConstraint.GTC, OrderType.LIMIT_ORDER ) - ) as SimpleOrder - //then - Assertions.assertEquals(orderBook.bidOrders.entriesList().size, 2) - Assertions.assertEquals(orderBook.bestBidOrder, order) - Assertions.assertEquals(bestBidOrder!!.better, order) - Assertions.assertEquals(order.worse, bestBidOrder) - Assertions.assertEquals(orderBook.bidOrders.get(order.price).lastOrder, order) - Assertions.assertEquals(orderBook.bidOrders.get(order.price).totalQuantity, 1) - Assertions.assertEquals(orderBook.bidOrders.get(order.price).ordersCount, 1) - } + ) + val bestBidOrder = orderBook.bestBidOrder + //when + val order: SimpleOrder = + orderBook.handleNewOrderCommand( + OrderCreateCommand( + UUID.randomUUID().toString(), + uuid, + pair, + 2, + 1, + OrderDirection.BID, + MatchConstraint.GTC, + OrderType.LIMIT_ORDER + ) + ) as SimpleOrder + //then + Assertions.assertEquals(orderBook.bidOrders.entriesList().size, 2) + Assertions.assertEquals(orderBook.bestBidOrder, order) + Assertions.assertEquals(bestBidOrder!!.better, order) + Assertions.assertEquals(order.worse, bestBidOrder) + Assertions.assertEquals(orderBook.bidOrders.get(order.price).lastOrder, order) + Assertions.assertEquals(orderBook.bidOrders.get(order.price).totalQuantity, 1) + Assertions.assertEquals(orderBook.bidOrders.get(order.price).ordersCount, 1) + } @Test - fun givenOrderBookWithBidOrders_whenGtcAskLimitOrderWithSamePriceCreated_thenInstantMatch() { + fun givenOrderBookWithBidOrders_whenGtcAskLimitOrderWithSamePriceCreated_thenInstantMatch(): Unit = runBlocking { //given - val orderBook = SimpleOrderBook(pair, false) + val orderBook = createOrderBook() orderBook.handleNewOrderCommand( OrderCreateCommand( UUID.randomUUID().toString(), @@ -201,9 +204,9 @@ class SimpleOrderBookUnitTest { } @Test - fun givenOrderBookWithBidOrders_whenGtcAskLimitOrderWithNotMatchPriceCreated_thenAddToQueue() { + fun givenOrderBookWithBidOrders_whenGtcAskLimitOrderWithNotMatchPriceCreated_thenAddToQueue(): Unit = runBlocking { //given - val orderBook = SimpleOrderBook(pair, false) + val orderBook = createOrderBook() orderBook.handleNewOrderCommand( OrderCreateCommand( UUID.randomUUID().toString(), @@ -250,70 +253,71 @@ class SimpleOrderBookUnitTest { } @Test - fun givenOrderBookWithBidAndAskOrders_whenGtcAskLimitOrderWithMatchPriceGreaterQuantityCreated_thenAddToQueue() { - //given - val orderBook = SimpleOrderBook(pair, false) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 2, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 1, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 3, - 1, - OrderDirection.ASK, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER + fun givenOrderBookWithBidAndAskOrders_whenGtcAskLimitOrderWithMatchPriceGreaterQuantityCreated_thenAddToQueue(): Unit = + runBlocking { + //given + val orderBook = createOrderBook() + orderBook.handleNewOrderCommand( + OrderCreateCommand( + UUID.randomUUID().toString(), + uuid, + pair, + 2, + 1, + OrderDirection.BID, + MatchConstraint.GTC, + OrderType.LIMIT_ORDER + ) ) - ) - //when - val order: SimpleOrder = orderBook.handleNewOrderCommand( OrderCreateCommand( UUID.randomUUID().toString(), uuid, pair, 1, + 1, + OrderDirection.BID, + MatchConstraint.GTC, + OrderType.LIMIT_ORDER + ) + ) + orderBook.handleNewOrderCommand( + OrderCreateCommand( + UUID.randomUUID().toString(), + uuid, + pair, 3, + 1, OrderDirection.ASK, MatchConstraint.GTC, OrderType.LIMIT_ORDER ) - ) as SimpleOrder - //then - Assertions.assertEquals(orderBook.bidOrders.entriesList().size, 0) - Assertions.assertEquals(orderBook.askOrders.entriesList().size, 2) - Assertions.assertNull(orderBook.bestBidOrder) - Assertions.assertEquals(orderBook.bestAskOrder, order) - } + ) + //when + val order: SimpleOrder = + orderBook.handleNewOrderCommand( + OrderCreateCommand( + UUID.randomUUID().toString(), + uuid, + pair, + 1, + 3, + OrderDirection.ASK, + MatchConstraint.GTC, + OrderType.LIMIT_ORDER + ) + ) as SimpleOrder + //then + Assertions.assertEquals(orderBook.bidOrders.entriesList().size, 0) + Assertions.assertEquals(orderBook.askOrders.entriesList().size, 2) + Assertions.assertNull(orderBook.bestBidOrder) + Assertions.assertEquals(orderBook.bestAskOrder, order) + } @Test - fun givenOrderBook_whenCancelBestBidOrder_thenBestBidOrderChange() { + fun givenOrderBook_whenCancelBestBidOrder_thenBestBidOrderChange(): Unit = runBlocking { //given - val orderBook = SimpleOrderBook(pair, false) + val orderBook = createOrderBook() val firstOrderId = UUID.randomUUID().toString() val secondOrderId = UUID.randomUUID().toString() @@ -349,9 +353,9 @@ class SimpleOrderBookUnitTest { } @Test - fun givenOrderBookWithMoreBids_whenCancelBestBidOrder_thenBestBidOrderChange() { + fun givenOrderBookWithMoreBids_whenCancelBestBidOrder_thenBestBidOrderChange(): Unit = runBlocking { //given - val orderBook = SimpleOrderBook(pair, false) + val orderBook = createOrderBook() val firstOrderId = UUID.randomUUID().toString() val secondOrderId = UUID.randomUUID().toString() @@ -399,9 +403,9 @@ class SimpleOrderBookUnitTest { } @Test - fun givenOrderBookWithMoreBids_whenCancelABidOrder_thenBestBidOrderNotChange() { + fun givenOrderBookWithMoreBids_whenCancelABidOrder_thenBestBidOrderNotChange(): Unit = runBlocking { //given - val orderBook = SimpleOrderBook(pair, false) + val orderBook = createOrderBook() val firstOrderId = UUID.randomUUID().toString() val secondOrderId = UUID.randomUUID().toString() @@ -450,148 +454,9 @@ class SimpleOrderBookUnitTest { @Test - fun givenOrderBookWithMoreBids_whenEditABidOrder_thenBestBidOrderChange() { - //given - val orderBook = SimpleOrderBook(pair, false) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 2, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - val secondOrder = orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 2, - 3, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 1, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - //when - val order = orderBook.handleEditCommand( - OrderEditCommand( - UUID.randomUUID().toString(), - uuid, - secondOrder!!.id()!!, - pair, - 3, - 2 - ) - ) - //then - Assertions.assertEquals(secondOrder.id(), order?.id()) - Assertions.assertEquals(orderBook.bestBidOrder, order) - Assertions.assertEquals(orderBook.bidOrders.entriesList().size, 3) - } - - @Test - fun givenOrderBookWithBidAndAskOrders_whenEditABidOrder_thenRefill() { - //given - val orderBook = SimpleOrderBook(pair, false) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 2, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - val secondBid = orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 1, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 3, - 1, - OrderDirection.ASK, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - //when - val order: SimpleOrder = orderBook.handleEditCommand( - OrderEditCommand( - UUID.randomUUID().toString(), - uuid, - secondBid!!.id()!!, - pair, - 3, - 3 - ) - ) as SimpleOrder - //then - Assertions.assertEquals(2, orderBook.bidOrders.entriesList().size) - Assertions.assertEquals(0, orderBook.askOrders.entriesList().size) - Assertions.assertEquals(orderBook.bestBidOrder, order) - Assertions.assertNull(orderBook.bestAskOrder) - } - - @Test - fun givenEmptyOrderBook_whenGtcBidMarketOrderCreated_thenRejected() { - //given - val orderBook = SimpleOrderBook(pair, false) - //when - - val order = orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 1, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.MARKET_ORDER - ) - ) - //then - Assertions.assertEquals(orderBook.bidOrders.entriesList().size, 0) - Assertions.assertNull(orderBook.bestBidOrder) - Assertions.assertNull(order) - } - - @Test - fun givenEmptyOrderBook_whenIocBidMarketOrderCreated_thenNoOrderCreated() { + fun givenEmptyOrderBook_whenGtcBidMarketOrderCreated_thenRejected(): Unit = runBlocking { //given - val orderBook = SimpleOrderBook(pair, false) + val orderBook = createOrderBook() //when val order = orderBook.handleNewOrderCommand( @@ -613,136 +478,138 @@ class SimpleOrderBookUnitTest { } @Test - fun givenOrderBookWithBidAndAskOrders_whenIocAskMarketOrderWithGreaterQuantityCreated_thenPartiallyFilled() { - //given - val orderBook = SimpleOrderBook(pair, false) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 2, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 1, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER + fun givenOrderBookWithBidAndAskOrders_whenIocAskMarketOrderWithGreaterQuantityCreated_thenPartiallyFilled(): Unit = + runBlocking { + //given + val orderBook = createOrderBook() + orderBook.handleNewOrderCommand( + OrderCreateCommand( + UUID.randomUUID().toString(), + uuid, + pair, + 2, + 1, + OrderDirection.BID, + MatchConstraint.GTC, + OrderType.LIMIT_ORDER + ) ) - ) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 3, - 1, - OrderDirection.ASK, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER + orderBook.handleNewOrderCommand( + OrderCreateCommand( + UUID.randomUUID().toString(), + uuid, + pair, + 1, + 1, + OrderDirection.BID, + MatchConstraint.GTC, + OrderType.LIMIT_ORDER + ) ) - ) - val bestAskOrder = orderBook.bestAskOrder - //when - val order: SimpleOrder = orderBook.handleNewOrderCommand( OrderCreateCommand( UUID.randomUUID().toString(), uuid, pair, - 0, 3, + 1, OrderDirection.ASK, - MatchConstraint.IOC, - OrderType.MARKET_ORDER + MatchConstraint.GTC, + OrderType.LIMIT_ORDER ) - ) as SimpleOrder - //then - Assertions.assertEquals(2, order.filledQuantity) - Assertions.assertEquals(orderBook.bidOrders.entriesList().size, 0) - Assertions.assertEquals(orderBook.askOrders.entriesList().size, 1) - Assertions.assertNull(orderBook.bestBidOrder) - Assertions.assertEquals(orderBook.bestAskOrder, bestAskOrder) - } + ) + val bestAskOrder = orderBook.bestAskOrder + //when + val order: SimpleOrder = + orderBook.handleNewOrderCommand( + OrderCreateCommand( + UUID.randomUUID().toString(), + uuid, + pair, + 0, + 3, + OrderDirection.ASK, + MatchConstraint.IOC, + OrderType.MARKET_ORDER + ) + ) as SimpleOrder + //then + Assertions.assertEquals(2, order.filledQuantity) + Assertions.assertEquals(orderBook.bidOrders.entriesList().size, 0) + Assertions.assertEquals(orderBook.askOrders.entriesList().size, 1) + Assertions.assertNull(orderBook.bestBidOrder) + Assertions.assertEquals(orderBook.bestAskOrder, bestAskOrder) + } @Test - fun givenOrderBookWithBidAndAskOrders_whenIocAskLimitOrderWithHigherPriceAndGreaterQuantityCreated_thenNotFilled() { - //given - val orderBook = SimpleOrderBook(pair, false) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 2, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER - ) - ) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 1, - 1, - OrderDirection.BID, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER + fun givenOrderBookWithBidAndAskOrders_whenIocAskLimitOrderWithHigherPriceAndGreaterQuantityCreated_thenNotFilled(): Unit = + runBlocking { + //given + val orderBook = createOrderBook() + orderBook.handleNewOrderCommand( + OrderCreateCommand( + UUID.randomUUID().toString(), + uuid, + pair, + 2, + 1, + OrderDirection.BID, + MatchConstraint.GTC, + OrderType.LIMIT_ORDER + ) ) - ) - orderBook.handleNewOrderCommand( - OrderCreateCommand( - UUID.randomUUID().toString(), - uuid, - pair, - 3, - 1, - OrderDirection.ASK, - MatchConstraint.GTC, - OrderType.LIMIT_ORDER + orderBook.handleNewOrderCommand( + OrderCreateCommand( + UUID.randomUUID().toString(), + uuid, + pair, + 1, + 1, + OrderDirection.BID, + MatchConstraint.GTC, + OrderType.LIMIT_ORDER + ) ) - ) - val bestAskOrder = orderBook.bestAskOrder - val bestBidOrder = orderBook.bestBidOrder - //when - val order: SimpleOrder = orderBook.handleNewOrderCommand( OrderCreateCommand( UUID.randomUUID().toString(), uuid, pair, 3, - 3, + 1, OrderDirection.ASK, - MatchConstraint.IOC, + MatchConstraint.GTC, OrderType.LIMIT_ORDER ) - ) as SimpleOrder - //then - Assertions.assertEquals(0, order.filledQuantity) - Assertions.assertEquals(orderBook.bidOrders.entriesList().size, 2) - Assertions.assertEquals(orderBook.askOrders.entriesList().size, 1) - Assertions.assertEquals(bestBidOrder, orderBook.bestBidOrder) - Assertions.assertEquals(bestAskOrder, orderBook.bestAskOrder) - } + ) + val bestAskOrder = orderBook.bestAskOrder + val bestBidOrder = orderBook.bestBidOrder + //when + val order: SimpleOrder = + orderBook.handleNewOrderCommand( + OrderCreateCommand( + UUID.randomUUID().toString(), + uuid, + pair, + 3, + 3, + OrderDirection.ASK, + MatchConstraint.IOC, + OrderType.LIMIT_ORDER + ) + ) as SimpleOrder + //then + Assertions.assertEquals(0, order.filledQuantity) + Assertions.assertEquals(orderBook.bidOrders.entriesList().size, 2) + Assertions.assertEquals(orderBook.askOrders.entriesList().size, 1) + Assertions.assertEquals(bestBidOrder, orderBook.bestBidOrder) + Assertions.assertEquals(bestAskOrder, orderBook.bestAskOrder) + } @Test - fun whenSample1SequenceOfOrdersOccurs_thenAllSuccess() { + fun whenSample1SequenceOfOrdersOccurs_thenAllSuccess(): Unit = runBlocking { - val orderBook = SimpleOrderBook(ETH_BTC_PAIR, false) + val orderBook = SimpleOrderBook(ETH_BTC_PAIR, false, CollectingOrderBookEventSink()) orderBook.handleNewOrderCommand( OrderCreateCommand( UUID.randomUUID().toString(), @@ -849,4 +716,7 @@ class SimpleOrderBookUnitTest { Assertions.assertNotNull(orderBook.bestBidOrder) Assertions.assertNotNull(orderBook.bestAskOrder) } + private fun createOrderBook(): SimpleOrderBook { + return SimpleOrderBook(pair, false, CollectingOrderBookEventSink()) + } } \ No newline at end of file diff --git a/matching-engine/matching-engine-ports/matching-engine-eventlistener-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/listener/consumer/OrderKafkaListener.kt b/matching-engine/matching-engine-ports/matching-engine-eventlistener-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/listener/consumer/OrderKafkaListener.kt index 16fd32a75..7c05acc72 100644 --- a/matching-engine/matching-engine-ports/matching-engine-eventlistener-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/listener/consumer/OrderKafkaListener.kt +++ b/matching-engine/matching-engine-ports/matching-engine-eventlistener-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/listener/consumer/OrderKafkaListener.kt @@ -5,6 +5,7 @@ import co.nilin.opex.matching.engine.ports.kafka.listener.spi.OrderRequestEventL import kotlinx.coroutines.runBlocking import org.apache.kafka.clients.consumer.ConsumerRecord import org.springframework.kafka.listener.MessageListener +import org.springframework.kafka.support.KafkaUtils import org.springframework.stereotype.Component @Component @@ -15,7 +16,7 @@ class OrderKafkaListener : MessageListener { override fun onMessage(data: ConsumerRecord) { orderListeners.forEach { tl -> runBlocking { - tl.onOrder(data.value(), data.partition(), data.offset(), data.timestamp()) + tl.onOrder(data.value(), data.partition(), data.offset(), data.timestamp(), data.topic(), KafkaUtils.getConsumerGroupId()) } } } diff --git a/matching-engine/matching-engine-ports/matching-engine-eventlistener-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/listener/spi/OrderRequestEventListener.kt b/matching-engine/matching-engine-ports/matching-engine-eventlistener-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/listener/spi/OrderRequestEventListener.kt index 32342e348..3e01989c1 100644 --- a/matching-engine/matching-engine-ports/matching-engine-eventlistener-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/listener/spi/OrderRequestEventListener.kt +++ b/matching-engine/matching-engine-ports/matching-engine-eventlistener-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/listener/spi/OrderRequestEventListener.kt @@ -4,5 +4,5 @@ import co.nilin.opex.matching.engine.core.inout.OrderRequestEvent interface OrderRequestEventListener { fun id(): String - suspend fun onOrder(order: OrderRequestEvent, partition: Int, offset: Long, timestamp: Long) + suspend fun onOrder(order: OrderRequestEvent, partition: Int, offset: Long, timestamp: Long, topic: String, consumerGroupId: String) } \ No newline at end of file diff --git a/matching-engine/matching-engine-ports/matching-engine-submitter-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/submitter/config/EventsKafkaConfig.kt b/matching-engine/matching-engine-ports/matching-engine-submitter-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/submitter/config/EventsKafkaConfig.kt index 667fceeac..89546f4b7 100644 --- a/matching-engine/matching-engine-ports/matching-engine-submitter-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/submitter/config/EventsKafkaConfig.kt +++ b/matching-engine/matching-engine-ports/matching-engine-submitter-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/submitter/config/EventsKafkaConfig.kt @@ -28,8 +28,14 @@ class EventsKafkaConfig { ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG to StringSerializer::class.java, ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG to JsonSerializer::class.java, ProducerConfig.ACKS_CONFIG to "all", + ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG to true, + ProducerConfig.RETRIES_CONFIG to 10, + ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG to 5000, + ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG to 5000, + ProducerConfig.TRANSACTIONAL_ID_CONFIG to "matching-engine-tx", JsonDeserializer.TRUSTED_PACKAGES to "co.nilin.opex.*", JsonDeserializer.TYPE_MAPPINGS to "orderBookUpdate:co.nilin.opex.matching.engine.core.inout.OrderBookUpdateEvent" + ) } diff --git a/matching-engine/matching-engine-ports/matching-engine-submitter-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/submitter/service/OrderBookTransitionSubmitter.kt b/matching-engine/matching-engine-ports/matching-engine-submitter-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/submitter/service/OrderBookTransitionSubmitter.kt new file mode 100644 index 000000000..19eb9d3e0 --- /dev/null +++ b/matching-engine/matching-engine-ports/matching-engine-submitter-kafka/src/main/kotlin/co/nilin/opex/matching/engine/ports/kafka/submitter/service/OrderBookTransitionSubmitter.kt @@ -0,0 +1,127 @@ +package co.nilin.opex.matching.engine.ports.kafka.submitter.service + +import co.nilin.opex.matching.engine.core.eventh.events.TradeEvent +import co.nilin.opex.matching.engine.core.inout.InputKafkaMetadata +import co.nilin.opex.matching.engine.core.model.PreparedCommandResult +import co.nilin.opex.matching.engine.core.spi.OrderBookTransitionPublisher +import org.apache.kafka.clients.consumer.OffsetAndMetadata +import org.apache.kafka.common.TopicPartition +import org.springframework.kafka.core.KafkaTemplate +import org.springframework.stereotype.Component +import org.springframework.util.concurrent.ListenableFuture +import java.util.concurrent.ExecutionException +import co.nilin.opex.matching.engine.core.model.Pair +@Component +class OrderBookTransitionSubmitter( + private val kafkaTemplate: KafkaTemplate +) : OrderBookTransitionPublisher { + override suspend fun publish( + prepared: PreparedCommandResult, + inputMetadata: InputKafkaMetadata + ) { + check(kafkaTemplate.isTransactional) { + "KafkaTemplate is not configured for transactions" + } + + kafkaTemplate.executeInTransaction { operations -> + val futures = + mutableListOf>() + + /* + * Publish every public event generated by this command. + */ + prepared.events.forEach { event -> + val key = pairKey(event.pair) + + futures += operations.send( + eventsTopic(event.pair), + key, + event + ) + + /* + * Preserve your current behavior: + * TradeEvent is sent to both the general events topic + * and the dedicated trades topic. + */ + if (event is TradeEvent) { + futures += operations.send( + tradesTopic(event.pair), + key, + event + ) + } + } + + /* + * Publish the durable order-book delta only when + * the command changed order-book state. + */ + prepared.stateTransition?.let { transition -> + futures += operations.send( + ORDER_BOOK_CHANGELOG_TOPIC, + pairKey(prepared.pair), + transition.snapshot + ) + } + + /* + * Wait for all Kafka sends. + * + * If one send fails, this throws and executeInTransaction() + * aborts the complete transaction. + */ + waitForAll(futures) + + /* + * Commit the consumed input record's next offset + * in the same Kafka transaction. + */ + val offsets = mapOf( + TopicPartition( + inputMetadata.topic, + inputMetadata.partition + ) to OffsetAndMetadata( + inputMetadata.nextOffset + ) + ) + + operations.sendOffsetsToTransaction( + offsets, + inputMetadata.consumerGroupId + ) + + Unit + } + } + + private fun waitForAll( + futures: Iterable> + ) { + try { + futures.forEach { future -> + future.get() + } + } catch (exception: ExecutionException) { + /* + * Expose the real Kafka send failure instead of only + * the Future wrapper exception. + */ + throw exception.cause ?: exception + } + } + + private fun eventsTopic(pair: Pair ): String = + "events_${pair.leftSideName}_${pair.rightSideName}" + + private fun tradesTopic(pair: Pair): String = + "trades_${pair.leftSideName}_${pair.rightSideName}" + + private fun pairKey(pair: Pair): String = + "${pair.leftSideName}_${pair.rightSideName}" + + companion object { + private const val ORDER_BOOK_CHANGELOG_TOPIC = + "orderbook_changelog" + } +} \ No newline at end of file diff --git a/matching-engine/pom.xml b/matching-engine/pom.xml index 010d84c41..e336f2caf 100644 --- a/matching-engine/pom.xml +++ b/matching-engine/pom.xml @@ -28,6 +28,10 @@ org.springframework.boot spring-boot-starter-test + + co.nilin.opex + common + diff --git a/wallet/wallet-app/src/main/kotlin/co/nilin/opex/wallet/app/controller/WalletOwnerController.kt b/wallet/wallet-app/src/main/kotlin/co/nilin/opex/wallet/app/controller/WalletOwnerController.kt index 571db4d57..ca74f3bb5 100644 --- a/wallet/wallet-app/src/main/kotlin/co/nilin/opex/wallet/app/controller/WalletOwnerController.kt +++ b/wallet/wallet-app/src/main/kotlin/co/nilin/opex/wallet/app/controller/WalletOwnerController.kt @@ -10,6 +10,7 @@ import co.nilin.opex.wallet.core.spi.WalletOwnerManager import io.swagger.annotations.ApiResponse import io.swagger.annotations.Example import io.swagger.annotations.ExampleProperty +import org.apache.kafka.common.message.ListOffsetsRequestData import org.springframework.core.env.Environment import org.springframework.web.bind.annotation.GetMapping import org.springframework.web.bind.annotation.PathVariable @@ -37,14 +38,20 @@ class WalletOwnerController( ) ) ) - suspend fun getAllWallets(@PathVariable uuid: String): List { + suspend fun getAllWallets(@PathVariable uuid: String): List? { + val externalIdentifier= currentUserProvider.getCurrentUser()?.identityId val owner = walletOwnerManager.findWalletOwner(uuid) ?: run { if (currentUserProvider.getCurrentUser()?.uuid.equals(uuid) && environment.activeProfiles.contains("otc")) walletOwnerManager.createWalletOwner( - uuid, currentUserProvider.getCurrentUser()?.mobile ?: "not set", "", currentUserProvider.getCurrentUser()?.identityId + uuid, + currentUserProvider.getCurrentUser()?.mobile ?: "not set", + "", + externalIdentifier ) - throw OpexError.WalletOwnerNotFound.exception() + return null } + if (owner.externalIdentifier != externalIdentifier) + walletOwnerManager.updateWalletOwnerExternalIdentifier(uuid, externalIdentifier) val wallets = walletManager.findWalletsByOwner(owner) return balanceParser.parse(wallets) } diff --git a/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/model/WalletOwner.kt b/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/model/WalletOwner.kt index cbe67a13a..c8cf73cb1 100644 --- a/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/model/WalletOwner.kt +++ b/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/model/WalletOwner.kt @@ -7,5 +7,6 @@ data class WalletOwner( val level: String, val isTradeAllowed: Boolean, val isWithdrawAllowed: Boolean, - val isDepositAllowed: Boolean + val isDepositAllowed: Boolean, + val externalIdentifier: String? = null ) \ No newline at end of file diff --git a/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/spi/WalletOwnerManager.kt b/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/spi/WalletOwnerManager.kt index 9595b0051..74d2eb2ea 100644 --- a/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/spi/WalletOwnerManager.kt +++ b/wallet/wallet-core/src/main/kotlin/co/nilin/opex/wallet/core/spi/WalletOwnerManager.kt @@ -11,7 +11,15 @@ interface WalletOwnerManager { suspend fun isWithdrawAllowed(owner: WalletOwner, amount: Amount): Boolean suspend fun findWalletOwner(uuid: String): WalletOwner? suspend fun findWalletOwnerByExternalIdentifier(externalIdentifier: String): WalletOwner? - suspend fun createWalletOwner(uuid: String, title: String, userLevel: String, externalIdentifier: String?=null): WalletOwner + suspend fun createWalletOwner( + uuid: String, + title: String, + userLevel: String, + externalIdentifier: String? = null + ): WalletOwner + suspend fun findAllWalletOwners(): List suspend fun updateWalletOwnerName(uuid: String, name: String, externalIdentifier: String?) -} + suspend fun updateWalletOwnerExternalIdentifier(uuid: String, externalIdentifier: String?) + +} \ No newline at end of file diff --git a/wallet/wallet-core/src/test/kotlin/co/nilin/opex/wallet/core/service/WithdrawServiceTest.kt b/wallet/wallet-core/src/test/kotlin/co/nilin/opex/wallet/core/service/WithdrawServiceTest.kt index 5ed7bf48a..890cd5eb3 100644 --- a/wallet/wallet-core/src/test/kotlin/co/nilin/opex/wallet/core/service/WithdrawServiceTest.kt +++ b/wallet/wallet-core/src/test/kotlin/co/nilin/opex/wallet/core/service/WithdrawServiceTest.kt @@ -64,7 +64,7 @@ class WithdrawServiceTest { private fun createCurrency() = CurrencyCommand("BTC", UUID.randomUUID().toString(), "Bitcoin", BigDecimal.valueOf(0.0001)) - private fun createOwner() = WalletOwner(1L, USER_UUID, "User", "registered", true, true, true) + private fun createOwner() = WalletOwner(1L, USER_UUID, "User", "registered", true, true, true ) private fun createSystemOwner() = WalletOwner(2L, SYSTEM_UUID, "System", "registered", true, true, true) diff --git a/wallet/wallet-ports/wallet-persister-postgres/src/main/kotlin/co/nilin/opex/wallet/ports/postgres/impl/WalletOwnerManagerImpl.kt b/wallet/wallet-ports/wallet-persister-postgres/src/main/kotlin/co/nilin/opex/wallet/ports/postgres/impl/WalletOwnerManagerImpl.kt index 0d7329e39..ec32081bc 100644 --- a/wallet/wallet-ports/wallet-persister-postgres/src/main/kotlin/co/nilin/opex/wallet/ports/postgres/impl/WalletOwnerManagerImpl.kt +++ b/wallet/wallet-ports/wallet-persister-postgres/src/main/kotlin/co/nilin/opex/wallet/ports/postgres/impl/WalletOwnerManagerImpl.kt @@ -161,7 +161,8 @@ class WalletOwnerManagerImpl( it.level, it.isTradeAllowed, it.isWithdrawAllowed, - it.isDepositAllowed + it.isDepositAllowed, + it.externalIdentifier ) } } @@ -177,6 +178,14 @@ class WalletOwnerManagerImpl( logger.warn("Wallet owner not found for UUID: $uuid") } } + override suspend fun updateWalletOwnerExternalIdentifier(uuid: String,externalIdentifier: String?) { + val owner = walletOwnerRepository.findByUuid(uuid).awaitFirstOrNull() + if (owner != null) { + owner.externalIdentifier = externalIdentifier + walletOwnerRepository.save(owner).awaitFirstOrNull() + ?: logger.warn("Failed to update wallet externalIdentifier for UUID: $uuid") + } + } override suspend fun findWalletOwnerByExternalIdentifier(externalIdentifier: String): WalletOwner? { return walletOwnerRepository.findByExternalIdentifier(externalIdentifier).awaitFirstOrNull()?.toPlainObject()