diff --git a/api/api-ports/api-proxy-rest/src/main/kotlin/co/nilin/opex/api/ports/proxy/impl/MarketDataProxyImpl.kt b/api/api-ports/api-proxy-rest/src/main/kotlin/co/nilin/opex/api/ports/proxy/impl/MarketDataProxyImpl.kt index 41782c441..3791ef458 100644 --- a/api/api-ports/api-proxy-rest/src/main/kotlin/co/nilin/opex/api/ports/proxy/impl/MarketDataProxyImpl.kt +++ b/api/api-ports/api-proxy-rest/src/main/kotlin/co/nilin/opex/api/ports/proxy/impl/MarketDataProxyImpl.kt @@ -56,7 +56,7 @@ class MarketDataProxyImpl(@Qualifier("generalWebClient") private val webClient: .onStatus({ t -> t.isError }, { it.createException() }) .bodyToMono() .awaitSingleOrNull() - ?: PriceChange(symbol, openTime = Date().time, closeTime = interval.getTime()) + ?: PriceChange(symbol, openTime = interval.getTime(), closeTime = Date().time) } } diff --git a/market/market-ports/market-persister-postgres/src/main/kotlin/co/nilin/opex/market/ports/postgres/dao/TradeRepository.kt b/market/market-ports/market-persister-postgres/src/main/kotlin/co/nilin/opex/market/ports/postgres/dao/TradeRepository.kt index 517db2d74..8c6cd0a5a 100644 --- a/market/market-ports/market-persister-postgres/src/main/kotlin/co/nilin/opex/market/ports/postgres/dao/TradeRepository.kt +++ b/market/market-ports/market-persister-postgres/src/main/kotlin/co/nilin/opex/market/ports/postgres/dao/TradeRepository.kt @@ -189,26 +189,21 @@ interface TradeRepository : ReactiveCrudRepository { select symbol, (select matched_price from last_trade where symbol=t.symbol) - (select matched_price from first_trade where symbol=t.symbol) as price_change, ((((select matched_price from last_trade where symbol=t.symbol) - (select matched_price from first_trade where symbol=t.symbol))/(select matched_price from first_trade where symbol=t.symbol))*100) as price_change_percent, - (sum(matched_quantity)/sum(matched_price)) as weighted_avg_price, + (sum(matched_price * matched_quantity)/nullif(sum(matched_quantity), 0)) as weighted_avg_price, (select matched_price from last_trade where symbol=t.symbol) as last_price, (select matched_quantity from last_trade where symbol=t.symbol) as last_qty, ( - select price from orders + select max(price) from orders inner join open_orders oo on orders.ouid = oo.ouid where create_date > :date and symbol=t.symbol and side='BID' - order by create_date desc limit 1 ) as bid_price, ( - select price from orders + select min(price) from orders inner join open_orders oo on orders.ouid = oo.ouid where create_date > :date and symbol=t.symbol and side='ASK' - order by create_date desc limit 1 ) as ask_price, ( - select price from orders - inner join open_orders oo on orders.ouid = oo.ouid - where create_date > :date and symbol=t.symbol - order by create_date desc limit 1 + select matched_price from first_trade where symbol=t.symbol ) as open_price, max(matched_price) as high_price, min(matched_price) as low_price, @@ -230,26 +225,21 @@ interface TradeRepository : ReactiveCrudRepository { select symbol, (select matched_price from last_trade) - (select matched_price from first_trade) as price_change, ((((select matched_price from last_trade) - (select matched_price from first_trade))/(select matched_price from first_trade))*100) as price_change_percent, - (sum(matched_quantity)/sum(matched_price)) as weighted_avg_price, + (sum(matched_price * matched_quantity)/nullif(sum(matched_quantity), 0)) as weighted_avg_price, (select matched_price from last_trade) as last_price, (select matched_quantity from last_trade) as last_qty, ( - select price from orders + select max(price) from orders inner join open_orders oo on orders.ouid = oo.ouid where create_date > :date and symbol=t.symbol and side='BID' - order by create_date desc limit 1 ) as bid_price, ( - select price from orders + select min(price) from orders inner join open_orders oo on orders.ouid = oo.ouid where create_date > :date and symbol=t.symbol and side='ASK' - order by create_date desc limit 1 ) as ask_price, ( - select price from orders - inner join open_orders oo on orders.ouid = oo.ouid - where create_date > :date and symbol=t.symbol - order by create_date desc limit 1 + select matched_price from first_trade ) as open_price, max(matched_price) as high_price, min(matched_price) as low_price, @@ -350,29 +340,35 @@ interface TradeRepository : ReactiveCrudRepository { :interval::INTERVAL ) ), + limited_intervals AS ( + SELECT * + FROM intervals + ORDER BY start_time DESC + LIMIT :limit + ), first_trade AS ( - SELECT DISTINCT ON (f.start_time) - f.start_time, - f.end_time, + SELECT DISTINCT ON (i.start_time) + i.start_time, + i.end_time, t.matched_price AS open_price - FROM intervals f + FROM limited_intervals i LEFT JOIN trades t - ON t.create_date >= f.start_time - AND t.create_date < f.end_time + ON t.create_date >= i.start_time + AND t.create_date < i.end_time AND t.symbol = :symbol - ORDER BY f.start_time, t.create_date + ORDER BY i.start_time, t.create_date ), last_trade AS ( - SELECT DISTINCT ON (f.start_time) - f.start_time, - f.end_time, + SELECT DISTINCT ON (i.start_time) + i.start_time, + i.end_time, t.matched_price AS close_price - FROM intervals f + FROM limited_intervals i LEFT JOIN trades t - ON t.create_date >= f.start_time - AND t.create_date < f.end_time + ON t.create_date >= i.start_time + AND t.create_date < i.end_time AND t.symbol = :symbol - ORDER BY f.start_time, t.create_date DESC + ORDER BY i.start_time, t.create_date DESC ), ohlcv AS ( SELECT @@ -384,7 +380,7 @@ interface TradeRepository : ReactiveCrudRepository { lt.close_price AS close, SUM(t.matched_quantity) AS volume, COUNT(t.id) AS trades - FROM intervals i + FROM limited_intervals i LEFT JOIN trades t ON t.create_date >= i.start_time AND t.create_date < i.end_time @@ -396,12 +392,7 @@ interface TradeRepository : ReactiveCrudRepository { GROUP BY i.start_time, i.end_time, ft.open_price, lt.close_price ) SELECT * - FROM ( - SELECT * - FROM ohlcv - ORDER BY open_time DESC - limit :limit - ) sub + FROM ohlcv ORDER BY open_time ASC """ ) diff --git a/market/market-ports/market-persister-postgres/src/main/kotlin/co/nilin/opex/market/ports/postgres/impl/MarketQueryHandlerImpl.kt b/market/market-ports/market-persister-postgres/src/main/kotlin/co/nilin/opex/market/ports/postgres/impl/MarketQueryHandlerImpl.kt index fe8673002..5bce4f092 100644 --- a/market/market-ports/market-persister-postgres/src/main/kotlin/co/nilin/opex/market/ports/postgres/impl/MarketQueryHandlerImpl.kt +++ b/market/market-ports/market-persister-postgres/src/main/kotlin/co/nilin/opex/market/ports/postgres/impl/MarketQueryHandlerImpl.kt @@ -21,6 +21,7 @@ import java.math.BigDecimal import java.time.Instant import java.time.LocalDateTime import java.time.ZoneId +import java.time.temporal.ChronoUnit import java.util.* @@ -35,19 +36,23 @@ class MarketQueryHandlerImpl( override suspend fun getTradeTickerData(interval: Interval): List { return redisCacheHelper.getOrElse("tradeTickerData:${interval.label}", 2.minutes()) { + val closeTime = Date().time + val openTime = interval.getTime() tradeRepository.tradeTicker(interval.getLocalDateTime()) .collectList() .awaitFirstOrElse { emptyList() } - .map { it.asPriceChangeResponse(Date().time, interval.getTime()) } + .map { it.asPriceChangeResponse(openTime, closeTime) } } } override suspend fun getTradeTickerDateBySymbol(symbol: String, interval: Interval): PriceChange? { val cacheId = "tradeTickerData:$symbol:${interval.label}" return redisCacheHelper.getOrElse(cacheId, 2.minutes()) { + val closeTime = Date().time + val openTime = interval.getTime() tradeRepository.tradeTickerBySymbol(symbol, interval.getLocalDateTime()) .awaitSingleOrNull() - ?.asPriceChangeResponse(Date().time, interval.getTime()) + ?.asPriceChangeResponse(openTime, closeTime) } } @@ -283,21 +288,22 @@ class MarketQueryHandlerImpl( endTime: Long?, limit: Int, ): List { - val st = if (startTime == null) - tradeRepository.findFirstByCreateDate().awaitSingleOrNull()?.createDate ?: LocalDateTime.now() + val intervalStep = parseIntervalStep(interval) + val latestTradeDate = if (startTime == null || endTime == null) + tradeRepository.findLastByCreateDate().awaitSingleOrNull()?.createDate else - with(Instant.ofEpochMilli(startTime)) { - LocalDateTime.ofInstant(this, ZoneId.systemDefault()) - } - - val et = if (endTime == null) - tradeRepository.findLastByCreateDate().awaitSingleOrNull()?.createDate ?: LocalDateTime.now() - else - with(Instant.ofEpochMilli(endTime)) { - LocalDateTime.ofInstant(this, ZoneId.systemDefault()) - } + null + val fallbackDate = latestTradeDate ?: LocalDateTime.now() + val startDate = startTime?.asLocalDateTime() ?: when { + endTime != null -> shiftByIntervals(endTime.asLocalDateTime(), intervalStep, -(limit - 1).toLong()) + else -> shiftByIntervals(fallbackDate, intervalStep, -(limit - 1).toLong()) + } + val endDate = endTime?.asLocalDateTime() ?: when { + startTime != null -> shiftByIntervals(startDate, intervalStep, (limit - 1).toLong()) + else -> fallbackDate + } - return tradeRepository.candleData(symbol, interval, st, et, limit) + return tradeRepository.candleData(symbol, interval, startDate, endDate, limit) .collectList() .awaitFirstOrElse { emptyList() } .map { @@ -457,6 +463,32 @@ class MarketQueryHandlerImpl( count ?: 0 ) + private fun Long.asLocalDateTime(): LocalDateTime = with(Instant.ofEpochMilli(this)) { + LocalDateTime.ofInstant(this, ZoneId.systemDefault()) + } + + private fun parseIntervalStep(interval: String): Pair { + val parts = interval.trim().split(Regex("\\s+"), limit = 2) + val amount = parts.firstOrNull()?.toLongOrNull() + ?: throw IllegalArgumentException("Invalid interval amount: $interval") + val unit = when (parts.getOrNull(1)?.uppercase(Locale.US)?.removeSuffix("S")) { + "MINUTE" -> ChronoUnit.MINUTES + "HOUR" -> ChronoUnit.HOURS + "DAY" -> ChronoUnit.DAYS + else -> throw IllegalArgumentException("Unsupported interval unit: $interval") + } + return amount to unit + } + + private fun shiftByIntervals( + dateTime: LocalDateTime, + intervalStep: Pair, + intervals: Long, + ): LocalDateTime { + val (amount, unit) = intervalStep + return dateTime.plus(intervals * amount, unit) + } + private fun Long.approximate(): Long { if (this < 10) return this diff --git a/market/market-ports/market-persister-postgres/src/main/resources/schema.sql b/market/market-ports/market-persister-postgres/src/main/resources/schema.sql index 294115909..2395840f1 100644 --- a/market/market-ports/market-persister-postgres/src/main/resources/schema.sql +++ b/market/market-ports/market-persister-postgres/src/main/resources/schema.sql @@ -71,6 +71,7 @@ CREATE TABLE IF NOT EXISTS trades ); CREATE INDEX IF NOT EXISTS idx_trades_symbol on trades (symbol); CREATE INDEX IF NOT EXISTS idx_trades_create_date on trades (create_date); +CREATE INDEX IF NOT EXISTS idx_trades_symbol_create_date on trades (symbol, create_date); ALTER TABLE trades ALTER COLUMN id TYPE BIGINT, diff --git a/market/market-ports/market-persister-postgres/src/test/kotlin/co/nilin/opex/market/ports/postgres/impl/MarketQueryHandlerTest.kt b/market/market-ports/market-persister-postgres/src/test/kotlin/co/nilin/opex/market/ports/postgres/impl/MarketQueryHandlerTest.kt index bb25377ae..b0d479467 100644 --- a/market/market-ports/market-persister-postgres/src/test/kotlin/co/nilin/opex/market/ports/postgres/impl/MarketQueryHandlerTest.kt +++ b/market/market-ports/market-persister-postgres/src/test/kotlin/co/nilin/opex/market/ports/postgres/impl/MarketQueryHandlerTest.kt @@ -1,5 +1,6 @@ package co.nilin.opex.market.ports.postgres.impl +import co.nilin.opex.common.utils.Interval import co.nilin.opex.market.core.inout.MarketTrade import co.nilin.opex.market.core.inout.Order import co.nilin.opex.market.core.inout.OrderDirection @@ -8,7 +9,10 @@ import co.nilin.opex.market.ports.postgres.dao.OrderRepository import co.nilin.opex.market.ports.postgres.dao.OrderStatusRepository import co.nilin.opex.market.ports.postgres.dao.TradeRepository import co.nilin.opex.market.ports.postgres.impl.sample.VALID +import co.nilin.opex.market.ports.postgres.model.CandleInfoData import co.nilin.opex.market.ports.postgres.model.LastPrice +import co.nilin.opex.market.ports.postgres.model.TradeModel +import co.nilin.opex.market.ports.postgres.model.TradeTickerData import co.nilin.opex.market.ports.postgres.util.RedisCacheHelper import io.mockk.coEvery import io.mockk.every @@ -18,6 +22,8 @@ import org.assertj.core.api.Assertions.assertThat import org.junit.jupiter.api.Test import reactor.core.publisher.Flux import reactor.core.publisher.Mono +import java.math.BigDecimal +import java.time.LocalDateTime class MarketQueryHandlerTest { private val orderRepository = mockk() @@ -132,5 +138,93 @@ class MarketQueryHandlerTest { assertThat(marketTradeResponses?.count()).isEqualTo(1) assertThat(marketTradeResponses?.first()).isEqualTo(VALID.MARKET_TRADE_RESPONSE) } -} + @Test + fun givenTickerData_whenTradeTickerRequested_thenTickerTimeWindowIsOrderedCorrectly(): Unit = runBlocking { + val tradeTickerData = TradeTickerData( + VALID.ETH_USDT, + BigDecimal.ONE, + BigDecimal.ONE, + BigDecimal.ONE, + BigDecimal.ONE, + BigDecimal.ONE, + BigDecimal.ONE, + BigDecimal.ONE, + BigDecimal.ONE, + BigDecimal.TEN, + BigDecimal.ONE, + BigDecimal.TEN, + 1L, + 2L, + 3L + ) + coEvery { + redisCacheHelper.getOrElse>( + eq("tradeTickerData:${Interval.TwentyFourHours.label}"), + any(), + any() + ) + } coAnswers { + thirdArg List>().invoke() + } + every { tradeRepository.tradeTicker(any()) } returns Flux.just(tradeTickerData) + + val priceChanges = marketQueryHandler.getTradeTickerData(Interval.TwentyFourHours) + + assertThat(priceChanges).hasSize(1) + assertThat(priceChanges.first().openTime).isLessThanOrEqualTo(priceChanges.first().closeTime) + } + + @Test + fun givenMissingCandleBounds_whenGetCandleInfo_thenOnlyLatestIntervalsAreRequested(): Unit = runBlocking { + val latestTradeDate = LocalDateTime.of(2024, 1, 1, 10, 15) + val expectedStartDate = latestTradeDate.minusHours(2) + val latestTrade = TradeModel( + 1L, + 1L, + VALID.ETH_USDT, + "ETH", + "USDT", + BigDecimal.TEN, + BigDecimal.ONE, + BigDecimal.TEN, + BigDecimal.TEN, + BigDecimal.ZERO, + BigDecimal.ZERO, + "ETH", + "USDT", + latestTradeDate, + "maker", + "taker", + "maker-user", + "taker-user", + latestTradeDate + ) + val candleInfo = CandleInfoData( + expectedStartDate, + expectedStartDate.plusHours(1), + BigDecimal.ONE, + BigDecimal.TWO, + BigDecimal.TWO, + BigDecimal.ONE, + BigDecimal.TEN, + 1 + ) + coEvery { tradeRepository.findLastByCreateDate() } returns Mono.just(latestTrade) + coEvery { + tradeRepository.candleData( + VALID.ETH_USDT, + "1 HOURS", + expectedStartDate, + latestTradeDate, + 3 + ) + } returns Flux.just(candleInfo) + + val candles = marketQueryHandler.getCandleInfo(VALID.ETH_USDT, "1 HOURS", null, null, 3) + + assertThat(candles).hasSize(1) + assertThat(candles.first().openTime).isEqualTo(expectedStartDate) + assertThat(candles.first().closeTime).isEqualTo(expectedStartDate.plusHours(1)) + } +}