From 9c28079a783a6ec44f36f7a8c88a9a0ba87a0729 Mon Sep 17 00:00:00 2001 From: Junyou Park Date: Mon, 10 Aug 2026 02:25:55 +0900 Subject: [PATCH] feat: merge canonical history into market bars --- .../api/marketdata/MarketBarController.java | 2 +- .../marketdata/MarketBarControllerTest.java | 9 +++- .../marketdata/MarketBarJsonCodec.java | 43 ++++++++++++++++++ .../marketdata/RedisMarketBarAdapter.java | 35 ++++++++++++--- .../marketdata/RedisMarketBarAdapterTest.java | 44 +++++++++++++++++++ 5 files changed, 125 insertions(+), 8 deletions(-) diff --git a/apps/backend-api/src/main/java/com/idea2strategy/backend/api/marketdata/MarketBarController.java b/apps/backend-api/src/main/java/com/idea2strategy/backend/api/marketdata/MarketBarController.java index 7eef8ce9..910d0af5 100644 --- a/apps/backend-api/src/main/java/com/idea2strategy/backend/api/marketdata/MarketBarController.java +++ b/apps/backend-api/src/main/java/com/idea2strategy/backend/api/marketdata/MarketBarController.java @@ -29,7 +29,7 @@ public MarketBarController(MarketBarService service) { public SnapshotResponse recent( @PathVariable UUID instrumentId, @RequestParam(defaultValue = "30m") String timeframe, - @RequestParam(defaultValue = "300") int limit) { + @RequestParam(defaultValue = "400") int limit) { MarketBarTimeframe parsedTimeframe = MarketBarTimeframe.parse(timeframe); MarketBarSnapshot snapshot = service.findRecentSnapshot(instrumentId, parsedTimeframe, limit); return new SnapshotResponse( diff --git a/apps/backend-api/src/test/java/com/idea2strategy/backend/api/marketdata/MarketBarControllerTest.java b/apps/backend-api/src/test/java/com/idea2strategy/backend/api/marketdata/MarketBarControllerTest.java index f96fab23..ed0262c8 100644 --- a/apps/backend-api/src/test/java/com/idea2strategy/backend/api/marketdata/MarketBarControllerTest.java +++ b/apps/backend-api/src/test/java/com/idea2strategy/backend/api/marketdata/MarketBarControllerTest.java @@ -1,5 +1,6 @@ package com.idea2strategy.backend.api.marketdata; +import static org.junit.jupiter.api.Assertions.assertEquals; import static org.springframework.test.web.servlet.request.MockMvcRequestBuilders.get; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.jsonPath; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @@ -13,6 +14,7 @@ import java.time.Instant; import java.util.List; import java.util.UUID; +import java.util.concurrent.atomic.AtomicInteger; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.springframework.test.web.servlet.MockMvc; @@ -21,13 +23,16 @@ class MarketBarControllerTest { private static final UUID AAPL_ID = UUID.fromString("70000000-0000-4000-8000-000000000001"); private MockMvc mvc; + private AtomicInteger requestedLimit; @BeforeEach void setUp() { + requestedLimit = new AtomicInteger(); MarketBarPort port = new MarketBarPort() { @Override public List findRecent( UUID instrumentId, MarketBarTimeframe timeframe, int limit) { + requestedLimit.set(limit); return List.of(new com.idea2strategy.backend.application.marketdata.MarketBar( "event-1", AAPL_ID, "ALPACA", "SIP", Instant.parse("2026-08-06T14:30:00Z"), 1, 0, @@ -48,13 +53,13 @@ public List findRece @Test void returnsChronologicalStrategyBarsForTheRequestedTimeframe() throws Exception { mvc.perform(get("/api/v1/market-data/instruments/{instrumentId}/bars", AAPL_ID) - .queryParam("timeframe", "4h") - .queryParam("limit", "300")) + .queryParam("timeframe", "4h")) .andExpect(status().isOk()) .andExpect(jsonPath("$.instrumentId").value(AAPL_ID.toString())) .andExpect(jsonPath("$.symbol").value("AAPL")) .andExpect(jsonPath("$.timeframe").value("4h")) .andExpect(jsonPath("$.bars[0].close").value(210.12)); + assertEquals(400, requestedLimit.get()); } @Test diff --git a/modules/backend-messaging/src/main/java/com/idea2strategy/backend/messaging/marketdata/MarketBarJsonCodec.java b/modules/backend-messaging/src/main/java/com/idea2strategy/backend/messaging/marketdata/MarketBarJsonCodec.java index d5135302..ae811722 100644 --- a/modules/backend-messaging/src/main/java/com/idea2strategy/backend/messaging/marketdata/MarketBarJsonCodec.java +++ b/modules/backend-messaging/src/main/java/com/idea2strategy/backend/messaging/marketdata/MarketBarJsonCodec.java @@ -5,6 +5,8 @@ import com.idea2strategy.backend.application.marketdata.MarketBar; import java.math.BigDecimal; import java.time.Instant; +import java.util.ArrayList; +import java.util.List; import java.util.UUID; final class MarketBarJsonCodec { @@ -34,6 +36,47 @@ MarketBar decode(String encoded, com.idea2strategy.backend.application.marketdat } } + List decodeHistory( + String encoded, + UUID instrumentId, + com.idea2strategy.backend.application.marketdata.MarketBarTimeframe timeframe) { + try { + JsonNode root = objectMapper.readTree(encoded); + if (root.path("schemaVersion").asInt() != 1) { + throw new IllegalArgumentException("Unsupported history schema version"); + } + if (!"all".equals(root.path("adjustment").asText())) { + throw new IllegalArgumentException("History bars must use adjustment=all"); + } + if (!timeframe.value().equals(root.path("timeframe").asText())) { + throw new IllegalArgumentException("History timeframe does not match the request"); + } + if (!instrumentId.toString().equals(root.path("instrumentId").asText())) { + throw new IllegalArgumentException("History instrument does not match the request"); + } + List bars = new ArrayList<>(); + for (JsonNode value : root.path("bars")) { + Instant occurredAt = Instant.parse(requiredText(value, "t")); + long sequence = Math.floorDiv( + occurredAt.getEpochSecond(), timeframe.minutes() * 60L); + bars.add(new MarketBar( + "history:" + instrumentId + ":" + timeframe.value() + ":" + occurredAt, + instrumentId, + "ALPACA", + "SIP", + occurredAt, + sequence, + 0, + decimal(value, "o"), decimal(value, "h"), + decimal(value, "l"), decimal(value, "c"), + decimal(value, "v"))); + } + return List.copyOf(bars); + } catch (Exception exception) { + throw new IllegalArgumentException("Invalid historical market bar payload", exception); + } + } + private static String requiredText(JsonNode node, String field) { String value = node.path(field).asText(); if (value.isBlank()) throw new IllegalArgumentException("Missing " + field); diff --git a/modules/backend-messaging/src/main/java/com/idea2strategy/backend/messaging/marketdata/RedisMarketBarAdapter.java b/modules/backend-messaging/src/main/java/com/idea2strategy/backend/messaging/marketdata/RedisMarketBarAdapter.java index 42b471e6..1dc7358a 100644 --- a/modules/backend-messaging/src/main/java/com/idea2strategy/backend/messaging/marketdata/RedisMarketBarAdapter.java +++ b/modules/backend-messaging/src/main/java/com/idea2strategy/backend/messaging/marketdata/RedisMarketBarAdapter.java @@ -14,6 +14,7 @@ import java.math.BigDecimal; import java.util.LinkedHashMap; import java.util.Map; +import java.util.TreeMap; public final class RedisMarketBarAdapter implements MarketBarPort, AutoCloseable { private final RedisClient client; @@ -49,11 +50,16 @@ public List findRecent(UUID instrumentId, MarketBarTimeframe timefram if (timeframe.displayOnly()) { return findRecentDisplayBars(instrumentId, timeframe, limit); } - List encoded = commands.sync().zrevrange(recentBarsKey(instrumentId, timeframe), 0, limit - 1L); - List bars = new ArrayList<>(encoded.size()); - encoded.forEach(value -> bars.add(codec.decode(value, timeframe))); - Collections.reverse(bars); - return List.copyOf(bars); + List encoded = commands.sync().zrevrange( + recentBarsKey(instrumentId, timeframe), 0, limit - 1L); + List live = new ArrayList<>(encoded.size()); + encoded.forEach(value -> live.add(codec.decode(value, timeframe))); + Collections.reverse(live); + String historical = commands.sync().get(historyBarsKey(instrumentId, timeframe)); + List history = historical == null + ? List.of() + : codec.decodeHistory(historical, instrumentId, timeframe); + return mergeCanonicalHistoryWithLive(history, live, limit); } private List findRecentDisplayBars( @@ -107,6 +113,25 @@ String recentBarsKey(UUID instrumentId, MarketBarTimeframe timeframe) { return "{" + keyPrefix + ":market}:bars:" + instrumentId + ":" + timeframe.value(); } + String historyBarsKey(UUID instrumentId, MarketBarTimeframe timeframe) { + return "{" + keyPrefix + ":market}:history:bars:" + + instrumentId + ":" + timeframe.value(); + } + + static List mergeCanonicalHistoryWithLive( + List history, List live, int limit) { + TreeMap merged = new TreeMap<>(); + history.forEach(bar -> merged.put(bar.occurredAt(), bar)); + Instant historicalCutoff = merged.isEmpty() ? null : merged.lastKey(); + live.stream() + .filter(bar -> historicalCutoff == null || bar.occurredAt().isAfter(historicalCutoff)) + .forEach(bar -> merged.put(bar.occurredAt(), bar)); + List ordered = new ArrayList<>(merged.values()); + return ordered.size() <= limit + ? List.copyOf(ordered) + : List.copyOf(ordered.subList(ordered.size() - limit, ordered.size())); + } + private static String requirePrefix(String value) { if (value == null || value.isBlank() || value.contains("{") || value.contains("}")) { throw new IllegalArgumentException("market-data.redis-key-prefix must be a plain non-empty value"); diff --git a/modules/backend-messaging/src/test/java/com/idea2strategy/backend/messaging/marketdata/RedisMarketBarAdapterTest.java b/modules/backend-messaging/src/test/java/com/idea2strategy/backend/messaging/marketdata/RedisMarketBarAdapterTest.java index 60e58400..5e667e65 100644 --- a/modules/backend-messaging/src/test/java/com/idea2strategy/backend/messaging/marketdata/RedisMarketBarAdapterTest.java +++ b/modules/backend-messaging/src/test/java/com/idea2strategy/backend/messaging/marketdata/RedisMarketBarAdapterTest.java @@ -34,6 +34,50 @@ void aggregatesPersistedDisplayMinutesIntoFiveMinuteOhlcv() { }); } + @Test + void mergesOnlyLiveBarsAfterTheCanonicalHistoryCutoff() { + Instant start = Instant.parse("2026-08-06T14:30:00Z"); + MarketBar historical = minute(start, 0); + MarketBar historicalLatest = minute(start.plusSeconds(1800), 1); + MarketBar staleBeforeHistory = new MarketBar( + "stale-before-history", INSTRUMENT, "ALPACA", "SIP", + start, 1, 1, + BigDecimal.valueOf(150), BigDecimal.valueOf(151), + BigDecimal.valueOf(149), BigDecimal.valueOf(150), BigDecimal.TEN); + MarketBar staleOverlap = new MarketBar( + "stale-overlap", INSTRUMENT, "ALPACA", "SIP", + start.plusSeconds(1800), 2, 1, + BigDecimal.valueOf(200), BigDecimal.valueOf(201), + BigDecimal.valueOf(199), BigDecimal.valueOf(200), BigDecimal.TEN); + MarketBar liveLatest = minute(start.plusSeconds(3600), 3); + + List result = RedisMarketBarAdapter.mergeCanonicalHistoryWithLive( + List.of(historical, historicalLatest), + List.of(staleBeforeHistory, staleOverlap, liveLatest), + 3); + + assertThat(result).extracting(MarketBar::eventId) + .containsExactly("minute-0", "minute-1", "minute-3"); + } + + @Test + void decodesCompactAdjustedHistoryProjection() { + String payload = """ + {"schemaVersion":1,"adjustment":"all","timeframe":"30m", + "instrumentId":"70000000-0000-4000-8000-000000000001","bars":[ + {"t":"2026-08-06T14:30:00Z","o":100,"h":102,"l":99,"c":101,"v":1200}]} + """; + + List result = new MarketBarJsonCodec().decodeHistory( + payload, INSTRUMENT, MarketBarTimeframe.THIRTY_MINUTES); + + assertThat(result).singleElement().satisfies(bar -> { + assertThat(bar.occurredAt()).isEqualTo("2026-08-06T14:30:00Z"); + assertThat(bar.close()).isEqualByComparingTo("101"); + assertThat(bar.provider()).isEqualTo("ALPACA"); + }); + } + private static MarketBar minute(Instant at, int index) { BigDecimal open = BigDecimal.valueOf(100 + index); return new MarketBar(