Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions common/src/main/kotlin/co/nilin/opex/common/OpexError.kt
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
@@ -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<String, OrderBook>()
private val orderBooks =
ConcurrentHashMap<String, SimpleOrderBook>()

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
}

}

}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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<String>

Expand Down Expand Up @@ -55,7 +57,7 @@ class AppConfig {

@Bean
fun orderListener(): OrderListener {
return OrderListener()
return OrderListener(orderCommandProcessor,recoveryManager)
}

@Autowired
Expand Down

This file was deleted.

Original file line number Diff line number Diff line change
@@ -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
)
}
Original file line number Diff line number Diff line change
@@ -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<String> {
return symbols.split(",").map { it.trim() }.map { it.uppercase() }
}
}
Original file line number Diff line number Diff line change
@@ -1,49 +1,53 @@
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)

override fun id(): String {
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}"
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading