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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
Original file line number Diff line number Diff line change
@@ -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;
Expand All @@ -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;
Expand All @@ -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<com.idea2strategy.backend.application.marketdata.MarketBar> 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,
Expand All @@ -48,13 +53,13 @@ public List<com.idea2strategy.backend.application.marketdata.MarketBar> 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -34,6 +36,47 @@ MarketBar decode(String encoded, com.idea2strategy.backend.application.marketdat
}
}

List<MarketBar> 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<MarketBar> 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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -49,11 +50,16 @@ public List<MarketBar> findRecent(UUID instrumentId, MarketBarTimeframe timefram
if (timeframe.displayOnly()) {
return findRecentDisplayBars(instrumentId, timeframe, limit);
}
List<String> encoded = commands.sync().zrevrange(recentBarsKey(instrumentId, timeframe), 0, limit - 1L);
List<MarketBar> bars = new ArrayList<>(encoded.size());
encoded.forEach(value -> bars.add(codec.decode(value, timeframe)));
Collections.reverse(bars);
return List.copyOf(bars);
List<String> encoded = commands.sync().zrevrange(
recentBarsKey(instrumentId, timeframe), 0, limit - 1L);
List<MarketBar> 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<MarketBar> history = historical == null
? List.of()
: codec.decodeHistory(historical, instrumentId, timeframe);
return mergeCanonicalHistoryWithLive(history, live, limit);
}

private List<MarketBar> findRecentDisplayBars(
Expand Down Expand Up @@ -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<MarketBar> mergeCanonicalHistoryWithLive(
List<MarketBar> history, List<MarketBar> live, int limit) {
TreeMap<Instant, MarketBar> 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<MarketBar> 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");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<MarketBar> 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<MarketBar> 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(
Expand Down