diff --git a/Makefile b/Makefile new file mode 100644 index 0000000..5facafe --- /dev/null +++ b/Makefile @@ -0,0 +1,3 @@ +dev-api-price: + docker-compose up -d + ./gradlew :app-api-price:bootRun diff --git a/app/app-api-price/build.gradle.kts b/app/app-api-price/build.gradle.kts index b5deb1b..c70b720 100644 --- a/app/app-api-price/build.gradle.kts +++ b/app/app-api-price/build.gradle.kts @@ -1,3 +1,5 @@ +import java.net.Socket + plugins { java id("org.springframework.boot") version "3.4.7" @@ -25,11 +27,19 @@ repositories { dependencies { implementation(project(":shared")) + implementation(project(":domain:domain-price")) + implementation("org.springframework.boot:spring-boot-starter-jdbc") + runtimeOnly("org.postgresql:postgresql") + implementation(project(":infra")) + implementation("org.springframework.boot:spring-boot-starter-web") implementation("org.springdoc:springdoc-openapi-starter-webmvc-ui:2.8.9") + compileOnly("org.projectlombok:lombok") - developmentOnly("org.springframework.boot:spring-boot-devtools") annotationProcessor("org.projectlombok:lombok") + + developmentOnly("org.springframework.boot:spring-boot-devtools") + testImplementation("org.springframework.boot:spring-boot-starter-test") testRuntimeOnly("org.junit.platform:junit-platform-launcher") } @@ -37,3 +47,32 @@ dependencies { tasks.withType { useJUnitPlatform() } + +fun waitForRedis(host: String, port: Int, timeoutSeconds: Int = 30) { + val deadline = System.currentTimeMillis() + timeoutSeconds * 1000 + while (System.currentTimeMillis() < deadline) { + try { + Socket(host, port).use { return } + } catch (_: Exception) { + Thread.sleep(500) + } + } + throw RuntimeException("Redis at $host:$port not available after $timeoutSeconds seconds.") +} + +tasks.register("waitForRedis") { + doLast { + println("⏳ Waiting for Redis to become available...") + waitForRedis("localhost", 6379) + println("✅ Redis is ready!") + } +} + +tasks.register("composeUp") { + workingDir = rootDir + commandLine = listOf("docker", "compose", "up", "-d") +} + +tasks.register("bootWithDocker") { + dependsOn("composeUp", "waitForRedis", "bootRun") +} diff --git a/app/app-api-price/src/main/java/com/polynomeer/app/api/price/AppApiPriceApplication.java b/app/app-api-price/src/main/java/com/polynomeer/app/api/price/AppApiPriceApplication.java index bcddba6..fa1f929 100644 --- a/app/app-api-price/src/main/java/com/polynomeer/app/api/price/AppApiPriceApplication.java +++ b/app/app-api-price/src/main/java/com/polynomeer/app/api/price/AppApiPriceApplication.java @@ -1,9 +1,18 @@ package com.polynomeer.app.api.price; +import com.polynomeer.domain.price.repository.PriceCacheProperties; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.context.annotation.ComponentScan; @SpringBootApplication +@ComponentScan(basePackages = { + "com.polynomeer.app.api.price", + "com.polynomeer.domain.price", + "com.polynomeer.infra", +}) +@EnableConfigurationProperties(PriceCacheProperties.class) public class AppApiPriceApplication { public static void main(String[] args) { diff --git a/app/app-api-price/src/main/java/com/polynomeer/app/api/price/controller/ChartController.java b/app/app-api-price/src/main/java/com/polynomeer/app/api/price/controller/ChartController.java new file mode 100644 index 0000000..d693a45 --- /dev/null +++ b/app/app-api-price/src/main/java/com/polynomeer/app/api/price/controller/ChartController.java @@ -0,0 +1,28 @@ +package com.polynomeer.app.api.price.controller; + +import com.polynomeer.domain.price.model.ChartPoint; +import com.polynomeer.domain.price.service.ChartQueryService; +import com.polynomeer.shared.common.dto.CommonResponse; +import lombok.RequiredArgsConstructor; +import org.springframework.web.bind.annotation.*; + +import java.time.ZonedDateTime; +import java.util.List; + +@RestController +@RequiredArgsConstructor +@RequestMapping("/api/v1/charts") +public class ChartController { + + private final ChartQueryService chartService; + + @GetMapping("/{tickerCode}") + public CommonResponse> getChart( + @PathVariable String tickerCode, + @RequestParam String interval, + @RequestParam ZonedDateTime from, + @RequestParam ZonedDateTime to) { + List response = chartService.getChart(tickerCode, interval, from, to); + return new CommonResponse<>("SUCCESS", response); + } +} diff --git a/app/app-api-price/src/main/java/com/polynomeer/app/api/price/controller/PriceController.java b/app/app-api-price/src/main/java/com/polynomeer/app/api/price/controller/PriceController.java index 265869f..24b960d 100644 --- a/app/app-api-price/src/main/java/com/polynomeer/app/api/price/controller/PriceController.java +++ b/app/app-api-price/src/main/java/com/polynomeer/app/api/price/controller/PriceController.java @@ -1,10 +1,11 @@ package com.polynomeer.app.api.price.controller; -import com.polynomeer.app.api.price.Price; import com.polynomeer.app.api.price.dto.PriceResponse; +import com.polynomeer.domain.price.service.PriceQueryService; +import com.polynomeer.domain.ticker.validation.TickerFormat; import com.polynomeer.shared.common.dto.CommonResponse; import com.polynomeer.shared.common.error.TickerErrorCode; -import com.polynomeer.shared.common.error.TickerNotFoundException; +import com.polynomeer.shared.common.error.TickerValidationException; import io.swagger.v3.oas.annotations.Operation; import io.swagger.v3.oas.annotations.Parameter; import io.swagger.v3.oas.annotations.responses.ApiResponse; @@ -12,7 +13,10 @@ import io.swagger.v3.oas.annotations.tags.Tag; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; -import org.springframework.web.bind.annotation.*; +import org.springframework.web.bind.annotation.GetMapping; +import org.springframework.web.bind.annotation.PathVariable; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RestController; @RestController @RequestMapping("/api/v1/quotes") @@ -21,6 +25,8 @@ @Tag(name = "Quotes", description = "Operations related to stock quotes") public class PriceController { + private final PriceQueryService priceQueryService; + @Operation(summary = "Get quote by ticker code", description = "Returns stock quote data for a given ticker code.") @ApiResponses(value = { @ApiResponse(responseCode = "200", description = "Successful response"), @@ -33,14 +39,17 @@ public CommonResponse getQuote( ) { log.info("getQuote tickerCode={}", tickerCode); - if (tickerCode.equals("error")) { - throw new TickerNotFoundException(TickerErrorCode.TICKER_NOT_FOUND); - } + validateTickerCode(tickerCode); + + var price = priceQueryService.getCurrentPrice(tickerCode); + var response = PriceResponse.from(price); + return new CommonResponse<>("SUCCESS", response); + } - return new CommonResponse<>( - "SUCCESS", - PriceResponse.from( - new Price("TEST", 1000L, 10L, 100.0, 10000L, null) - )); + private void validateTickerCode(String tickerCode) { + if (!TickerFormat.isValid(tickerCode)) { + throw new TickerValidationException(TickerErrorCode.TICKER_INVALID); + } } + } diff --git a/app/app-api-price/src/main/java/com/polynomeer/app/api/price/dto/PriceResponse.java b/app/app-api-price/src/main/java/com/polynomeer/app/api/price/dto/PriceResponse.java index b247ff4..976b2ea 100644 --- a/app/app-api-price/src/main/java/com/polynomeer/app/api/price/dto/PriceResponse.java +++ b/app/app-api-price/src/main/java/com/polynomeer/app/api/price/dto/PriceResponse.java @@ -1,6 +1,6 @@ package com.polynomeer.app.api.price.dto; -import com.polynomeer.app.api.price.Price; +import com.polynomeer.domain.price.model.Price; import java.time.ZonedDateTime; diff --git a/app/app-api-price/src/main/resources/application-local.yml b/app/app-api-price/src/main/resources/application-local.yml new file mode 100644 index 0000000..bb07977 --- /dev/null +++ b/app/app-api-price/src/main/resources/application-local.yml @@ -0,0 +1,30 @@ +spring: + config: + activate: + on-profile: local + + application: + name: app-api-price + + data: + redis: + host: localhost + port: 6379 + + datasource: + url: jdbc:postgresql://localhost:5432/romanticker + username: romanticker + password: romanticker + driver-class-name: org.postgresql.Driver + sql: + init: + mode: always + +logging: + pattern: + console: "%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{36} - [%X{traceId}] %msg%n" + + level: + com.polynomeer.infra.redis: DEBUG + com.polynomeer.infra.timescaledb: DEBUG + com.polynomeer.domain.price: DEBUG diff --git a/app/app-api-price/src/main/resources/application.properties b/app/app-api-price/src/main/resources/application.properties deleted file mode 100644 index 9e61395..0000000 --- a/app/app-api-price/src/main/resources/application.properties +++ /dev/null @@ -1,2 +0,0 @@ -spring.application.name=app-api-price -logging.pattern.console=%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{36} - [%X{traceId}] %msg%n diff --git a/app/app-api-price/src/main/resources/application.yml b/app/app-api-price/src/main/resources/application.yml new file mode 100644 index 0000000..83ff39d --- /dev/null +++ b/app/app-api-price/src/main/resources/application.yml @@ -0,0 +1,10 @@ +spring: + application: + name: app-api-price + + profiles: + active: local + +logging: + pattern: + console: "%d{yyyy-MM-dd HH:mm:ss.SSS} [%thread] %-5level %logger{36} - [%X{traceId}] %msg%n" diff --git a/app/app-api-price/src/main/resources/data.sql b/app/app-api-price/src/main/resources/data.sql new file mode 100644 index 0000000..6721d62 --- /dev/null +++ b/app/app-api-price/src/main/resources/data.sql @@ -0,0 +1,10 @@ +-- AAPL 1분 단위 시세 예시 +INSERT INTO price_history + (ticker_code, price, volume, "timestamp", exchange, currency, source) +VALUES ('AAPL', 19400, 100000, '2024-01-01T09:00:00Z', 'NASDAQ', 'USD', 'yahoo'), + ('AAPL', 19420, 120000, '2024-01-01T09:01:00Z', 'NASDAQ', 'USD', 'yahoo'), + ('AAPL', 19450, 150000, '2024-01-01T09:02:00Z', 'NASDAQ', 'USD', 'yahoo'), + ('AAPL', 19380, 130000, '2024-01-01T09:03:00Z', 'NASDAQ', 'USD', 'yahoo'), + ('AAPL', 19410, 160000, '2024-01-01T09:04:00Z', 'NASDAQ', 'USD', + 'yahoo') +ON CONFLICT (ticker_code, "timestamp") DO NOTHING; diff --git a/app/app-api-price/src/main/resources/schema.sql b/app/app-api-price/src/main/resources/schema.sql new file mode 100644 index 0000000..a922963 --- /dev/null +++ b/app/app-api-price/src/main/resources/schema.sql @@ -0,0 +1,25 @@ +CREATE EXTENSION IF NOT EXISTS timescaledb; + +CREATE TABLE IF NOT EXISTS price_history +( + ticker_code VARCHAR(20) NOT NULL, + "timestamp" TIMESTAMPTZ NOT NULL, + price BIGINT NOT NULL, + volume BIGINT NOT NULL, + exchange VARCHAR(20), + currency VARCHAR(10), + source VARCHAR(20) +); + +ALTER TABLE price_history + DROP CONSTRAINT IF EXISTS price_history_pkey; +DROP INDEX IF EXISTS ux_price_history_id; +DROP INDEX IF EXISTS ux_price_history_ticker_only; + +ALTER TABLE price_history + ADD CONSTRAINT price_history_pkey PRIMARY KEY (ticker_code, "timestamp"); + +SELECT create_hypertable('price_history', 'timestamp', if_not_exists => TRUE); + +CREATE INDEX IF NOT EXISTS ix_price_history_ticker_ts_desc + ON price_history (ticker_code, "timestamp" DESC); diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..e84ca9e --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,21 @@ +version: '3.8' +services: + redis: + image: redis:7 + ports: + - "6379:6379" + + timescaledb: + image: timescale/timescaledb:latest-pg15 + container_name: timescaledb + ports: + - "5432:5432" + environment: + POSTGRES_USER: romanticker + POSTGRES_PASSWORD: romanticker + POSTGRES_DB: romanticker + volumes: + - timescale_data:/var/lib/postgresql/data + +volumes: + timescale_data: diff --git a/domain/domain-price/build.gradle.kts b/domain/domain-price/build.gradle.kts index 16b058c..2eec2e7 100644 --- a/domain/domain-price/build.gradle.kts +++ b/domain/domain-price/build.gradle.kts @@ -18,7 +18,10 @@ repositories { } dependencies { + implementation(project(":shared")) implementation("org.springframework.boot:spring-boot-starter") + compileOnly("org.projectlombok:lombok") + annotationProcessor("org.projectlombok:lombok") testImplementation("org.springframework.boot:spring-boot-starter-test") testRuntimeOnly("org.junit.platform:junit-platform-launcher") } diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/.gitkeep b/domain/domain-price/src/main/java/com/polynomeer/domain/price/.gitkeep deleted file mode 100644 index e69de29..0000000 diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/model/ChartPoint.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/model/ChartPoint.java new file mode 100644 index 0000000..db7d457 --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/model/ChartPoint.java @@ -0,0 +1,13 @@ +package com.polynomeer.domain.price.model; + +import java.time.ZonedDateTime; + +public record ChartPoint( + ZonedDateTime timestamp, + long open, + long high, + long low, + long close, + long volume +) { +} diff --git a/app/app-api-price/src/main/java/com/polynomeer/app/api/price/Price.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/model/Price.java similarity index 82% rename from app/app-api-price/src/main/java/com/polynomeer/app/api/price/Price.java rename to domain/domain-price/src/main/java/com/polynomeer/domain/price/model/Price.java index 6033a14..0cf0aa7 100644 --- a/app/app-api-price/src/main/java/com/polynomeer/app/api/price/Price.java +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/model/Price.java @@ -1,4 +1,4 @@ -package com.polynomeer.app.api.price; +package com.polynomeer.domain.price.model; import java.time.ZonedDateTime; diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/BackoffStrategy.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/BackoffStrategy.java new file mode 100644 index 0000000..cef81b0 --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/BackoffStrategy.java @@ -0,0 +1,5 @@ +package com.polynomeer.domain.price.repository; + +public interface BackoffStrategy { + void pause(); +} \ No newline at end of file diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/CachePriceRepository.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/CachePriceRepository.java new file mode 100644 index 0000000..0aab8b9 --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/CachePriceRepository.java @@ -0,0 +1,13 @@ +package com.polynomeer.domain.price.repository; + +import com.polynomeer.domain.price.model.Price; + +import java.util.Optional; + +public interface CachePriceRepository { + Optional find(String tickerCode); + + void save(String tickerCode, Price latestFromDb); + + boolean saveIfAbsent(String tickerCode, Price price); +} diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/ExternalPriceClient.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/ExternalPriceClient.java new file mode 100644 index 0000000..ab6274d --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/ExternalPriceClient.java @@ -0,0 +1,7 @@ +package com.polynomeer.domain.price.repository; + +import com.polynomeer.domain.price.model.Price; + +public interface ExternalPriceClient { + Price fetchPrice(String tickerCode); +} diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/FixedBackoff.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/FixedBackoff.java new file mode 100644 index 0000000..fc85781 --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/FixedBackoff.java @@ -0,0 +1,20 @@ +package com.polynomeer.domain.price.repository; + +import java.time.Duration; + +public class FixedBackoff implements BackoffStrategy { + private final Duration d; + + public FixedBackoff(Duration d) { + this.d = d; + } + + @Override + public void pause() { + try { + Thread.sleep(d.toMillis()); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } +} \ No newline at end of file diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/PriceCacheConfig.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/PriceCacheConfig.java new file mode 100644 index 0000000..13ceb8b --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/PriceCacheConfig.java @@ -0,0 +1,31 @@ +package com.polynomeer.domain.price.repository; + +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Profile; + +import java.time.Duration; +import java.util.concurrent.Executor; +import java.util.concurrent.Executors; + +@Configuration +public class PriceCacheConfig { + + @Bean + public Executor priceQueryExecutor() { + return Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors()); + } + + @Bean + @Profile("!test") + public BackoffStrategy fixedBackoff() { + return new FixedBackoff(Duration.ofMillis(30)); + } + + @Bean + @Profile("test") + public BackoffStrategy noOpBackoff() { + return () -> { + }; + } +} diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/PriceCacheProperties.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/PriceCacheProperties.java new file mode 100644 index 0000000..fbefb46 --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/PriceCacheProperties.java @@ -0,0 +1,20 @@ +package com.polynomeer.domain.price.repository; + +import lombok.Getter; +import lombok.ToString; +import org.springframework.boot.context.properties.ConfigurationProperties; + +@ConfigurationProperties(prefix = "price.cache") +@Getter +@ToString +public class PriceCacheProperties { + private final int retries; + + public PriceCacheProperties() { + this.retries = 2; + } + + public PriceCacheProperties(int retries) { + this.retries = retries; + } +} diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/TimeSeriesChartRepository.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/TimeSeriesChartRepository.java new file mode 100644 index 0000000..7b8e9e0 --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/TimeSeriesChartRepository.java @@ -0,0 +1,10 @@ +package com.polynomeer.domain.price.repository; + +import com.polynomeer.domain.price.model.ChartPoint; + +import java.time.ZonedDateTime; +import java.util.List; + +public interface TimeSeriesChartRepository { + List findChart(String tickerCode, String interval, ZonedDateTime from, ZonedDateTime to); +} diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/TimeSeriesPriceRepository.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/TimeSeriesPriceRepository.java new file mode 100644 index 0000000..fda109e --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/repository/TimeSeriesPriceRepository.java @@ -0,0 +1,9 @@ +package com.polynomeer.domain.price.repository; + +import com.polynomeer.domain.price.model.Price; + +import java.util.Optional; + +public interface TimeSeriesPriceRepository { + Optional findLatest(String tickerCode); +} diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/ChartQueryService.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/ChartQueryService.java new file mode 100644 index 0000000..c95534c --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/ChartQueryService.java @@ -0,0 +1,10 @@ +package com.polynomeer.domain.price.service; + +import com.polynomeer.domain.price.model.ChartPoint; + +import java.time.ZonedDateTime; +import java.util.List; + +public interface ChartQueryService { + List getChart(String tickerCode, String interval, ZonedDateTime from, ZonedDateTime to); +} diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/ChartQueryServiceImpl.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/ChartQueryServiceImpl.java new file mode 100644 index 0000000..32bac33 --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/ChartQueryServiceImpl.java @@ -0,0 +1,21 @@ +package com.polynomeer.domain.price.service; + +import com.polynomeer.domain.price.model.ChartPoint; +import com.polynomeer.domain.price.repository.TimeSeriesChartRepository; +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Service; + +import java.time.ZonedDateTime; +import java.util.List; + +@Service +@RequiredArgsConstructor +public class ChartQueryServiceImpl implements ChartQueryService { + + private final TimeSeriesChartRepository chartRepository; + + @Override + public List getChart(String tickerCode, String interval, ZonedDateTime from, ZonedDateTime to) { + return chartRepository.findChart(tickerCode, interval, from, to); + } +} diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/PriceQueryService.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/PriceQueryService.java new file mode 100644 index 0000000..5a2ce82 --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/PriceQueryService.java @@ -0,0 +1,7 @@ +package com.polynomeer.domain.price.service; + +import com.polynomeer.domain.price.model.Price; + +public interface PriceQueryService { + Price getCurrentPrice(String tickerCode); +} diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/PriceQueryServiceImpl.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/PriceQueryServiceImpl.java new file mode 100644 index 0000000..63771ed --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/PriceQueryServiceImpl.java @@ -0,0 +1,63 @@ +package com.polynomeer.domain.price.service; + +import com.polynomeer.domain.price.model.Price; +import com.polynomeer.domain.price.repository.BackoffStrategy; +import com.polynomeer.domain.price.repository.CachePriceRepository; +import com.polynomeer.domain.price.repository.PriceCacheProperties; +import com.polynomeer.domain.price.repository.TimeSeriesPriceRepository; +import com.polynomeer.shared.common.error.PriceErrorCode; +import com.polynomeer.shared.common.error.PriceNotFoundException; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.stereotype.Service; + +import java.util.Optional; +import java.util.concurrent.Executor; + +@Slf4j +@Service +@RequiredArgsConstructor +public class PriceQueryServiceImpl implements PriceQueryService { + + private final CachePriceRepository redisRepository; + private final TimeSeriesPriceRepository timeSeriesRepository; + private final SingleFlightExecutor singleFlight; + @Qualifier("priceQueryExecutor") + private final Executor executor; + private final BackoffStrategy backoff; + private final PriceCacheProperties props; + + @Override + public Price getCurrentPrice(String tickerCode) { + return singleFlight.execute(tickerCode, () -> loadOnce(tickerCode), executor); + } + + private Price loadOnce(String code) { + return redisRepository.find(code).or(() -> { + log.debug("Cache miss for {}", code); + for (int i = 0; i < props.getRetries(); i++) { + backoff.pause(); + var again = redisRepository.find(code); + if (again.isPresent()) { + log.debug("Cache filled by peer for {}", code); + return again; + } + } + return Optional.empty(); + }).or(() -> { + var latest = timeSeriesRepository.findLatest(code) + .orElseThrow(() -> new PriceNotFoundException(PriceErrorCode.PRICE_NOT_FOUND)); + + var existing = redisRepository.find(code); + if (existing.isPresent() && existing.get().equals(latest)) { + log.debug("Skip redis write: identical value for {}", code); + return Optional.of(latest); + } + + boolean wrote = redisRepository.saveIfAbsent(code, latest); + log.debug("Redis write {}", wrote ? "SET(NX)" : "SKIPPED"); + return Optional.of(latest); + }).orElseThrow(() -> new PriceNotFoundException(PriceErrorCode.PRICE_NOT_FOUND)); + } +} diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/SingleFlightExecutor.java b/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/SingleFlightExecutor.java new file mode 100644 index 0000000..4c63d58 --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/price/service/SingleFlightExecutor.java @@ -0,0 +1,30 @@ +package com.polynomeer.domain.price.service; + +import org.springframework.stereotype.Component; + +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executor; +import java.util.function.Supplier; + +@Component +public class SingleFlightExecutor { + + private final ConcurrentHashMap> inFlight = new ConcurrentHashMap<>(); + + public V execute(K key, Supplier supplier, Executor executor) { + CompletableFuture future = inFlight.computeIfAbsent( + key, k -> CompletableFuture.supplyAsync(supplier, executor) + ); + try { + return future.join(); + } catch (CompletionException e) { + Throwable cause = e.getCause(); + if (cause instanceof RuntimeException re) throw re; + throw new RuntimeException(cause); + } finally { + inFlight.remove(key); + } + } +} diff --git a/domain/domain-price/src/main/java/com/polynomeer/domain/ticker/validation/TickerFormat.java b/domain/domain-price/src/main/java/com/polynomeer/domain/ticker/validation/TickerFormat.java new file mode 100644 index 0000000..a4efc2e --- /dev/null +++ b/domain/domain-price/src/main/java/com/polynomeer/domain/ticker/validation/TickerFormat.java @@ -0,0 +1,15 @@ +package com.polynomeer.domain.ticker.validation; + +import java.util.regex.Pattern; + +public final class TickerFormat { + + private TickerFormat() { + } + + public static final Pattern TICKER_PATTERN = Pattern.compile("^[A-Z]{1,5}([.-][A-Z0-9]{1,4})?$"); + + public static boolean isValid(String ticker) { + return ticker != null && TICKER_PATTERN.matcher(ticker).matches(); + } +} diff --git a/domain/domain-price/src/test/java/com/polynomeer/domain/price/.gitkeep b/domain/domain-price/src/test/java/com/polynomeer/domain/price/.gitkeep deleted file mode 100644 index e69de29..0000000 diff --git a/domain/domain-price/src/test/java/com/polynomeer/domain/price/service/FakeRepositories.java b/domain/domain-price/src/test/java/com/polynomeer/domain/price/service/FakeRepositories.java new file mode 100644 index 0000000..53815e4 --- /dev/null +++ b/domain/domain-price/src/test/java/com/polynomeer/domain/price/service/FakeRepositories.java @@ -0,0 +1,76 @@ +package com.polynomeer.domain.price.service; + +import com.polynomeer.domain.price.model.Price; +import com.polynomeer.domain.price.repository.BackoffStrategy; +import com.polynomeer.domain.price.repository.CachePriceRepository; +import com.polynomeer.domain.price.repository.TimeSeriesPriceRepository; + +import java.time.ZonedDateTime; +import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicInteger; + +public class FakeRepositories { + + static class FakeTimeSeriesRepository implements TimeSeriesPriceRepository { + final AtomicInteger reads = new AtomicInteger(); + final Price fixed; + final long delayMs; + + FakeTimeSeriesRepository(String code, double value, long delayMs) { + this.fixed = new Price(code, (long) (value * 100), 0, 0.0, 0, ZonedDateTime.now()); + this.delayMs = delayMs; + } + + @Override + public Optional findLatest(String tickerCode) { + reads.incrementAndGet(); + if (delayMs > 0) { + try { + Thread.sleep(delayMs); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + } + return Optional.of(fixed); + } + } + + static class NoopBackoff implements BackoffStrategy { + @Override + public void pause() { + } + } + + + static class FakeRedisRepository implements CachePriceRepository { + private final ConcurrentHashMap map = new ConcurrentHashMap<>(); + private final AtomicInteger writes = new AtomicInteger(); + + @Override + public Optional find(String tickerCode) { + return Optional.ofNullable(map.get(tickerCode)); + } + + @Override + public void save(String tickerCode, Price price) { + map.put(tickerCode, price); + writes.incrementAndGet(); + } + + @Override + public boolean saveIfAbsent(String tickerCode, Price price) { + Price prev = map.putIfAbsent(tickerCode, price); + if (prev == null) { + writes.incrementAndGet(); + return true; + } + return false; + } + + int writeCount() { + return writes.get(); + } + } + +} diff --git a/domain/domain-price/src/test/java/com/polynomeer/domain/price/service/PriceQueryServiceImplTest.java b/domain/domain-price/src/test/java/com/polynomeer/domain/price/service/PriceQueryServiceImplTest.java new file mode 100644 index 0000000..b7f7ad8 --- /dev/null +++ b/domain/domain-price/src/test/java/com/polynomeer/domain/price/service/PriceQueryServiceImplTest.java @@ -0,0 +1,186 @@ +package com.polynomeer.domain.price.service; + +import com.polynomeer.domain.price.model.Price; +import com.polynomeer.domain.price.repository.BackoffStrategy; +import com.polynomeer.domain.price.repository.CachePriceRepository; +import com.polynomeer.domain.price.repository.PriceCacheProperties; +import com.polynomeer.domain.price.repository.TimeSeriesPriceRepository; +import com.polynomeer.shared.common.error.PriceErrorCode; +import com.polynomeer.shared.common.error.PriceNotFoundException; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import java.time.ZonedDateTime; +import java.util.Optional; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.Executor; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; + +import static com.polynomeer.domain.price.service.FakeRepositories.FakeRedisRepository; +import static com.polynomeer.domain.price.service.FakeRepositories.FakeTimeSeriesRepository; +import static org.junit.jupiter.api.Assertions.*; +import static org.mockito.Mockito.*; + +class PriceQueryServiceImplTest { + + private CachePriceRepository cacheRepo; + private TimeSeriesPriceRepository dbRepo; + private PriceQueryServiceImpl priceService; + + private final Executor executor = Runnable::run; + private final BackoffStrategy noBackoff = () -> { + }; + private final PriceCacheProperties cacheProps = new PriceCacheProperties(); + + private final Price dummyPrice = + new Price("AAPL", 19500, 200, 1.03, 10_000_000L, ZonedDateTime.now()); + + @BeforeEach + void setUp() { + cacheRepo = mock(CachePriceRepository.class); + dbRepo = mock(TimeSeriesPriceRepository.class); + SingleFlightExecutor singleFlight = new SingleFlightExecutor<>(); + + priceService = new PriceQueryServiceImpl( + cacheRepo, dbRepo, singleFlight, executor, noBackoff, cacheProps + ); + } + + @Test + @DisplayName("캐시에 가격 정보가 있으면 DB를 조회하지 않고 바로 반환한다") + void shouldReturnPriceFromCacheIfExists() { + when(cacheRepo.find("AAPL")).thenReturn(Optional.of(dummyPrice)); + + Price result = priceService.getCurrentPrice("AAPL"); + + assertEquals(dummyPrice, result); + verify(cacheRepo, atLeastOnce()).find("AAPL"); + verifyNoInteractions(dbRepo); + verify(cacheRepo, never()).saveIfAbsent(anyString(), any()); + verify(cacheRepo, never()).save(anyString(), any()); + } + + @Test + @DisplayName("캐시에 없고 DB에 가격 정보가 있으면 DB에서 가져오고, 캐시에 NX로 저장한다") + void shouldReturnPriceFromDBIfCacheMiss() { + when(cacheRepo.find("AAPL")) + .thenReturn(Optional.empty()) + .thenReturn(Optional.empty()); + when(dbRepo.findLatest("AAPL")).thenReturn(Optional.of(dummyPrice)); + when(cacheRepo.saveIfAbsent("AAPL", dummyPrice)).thenReturn(true); + + Price result = priceService.getCurrentPrice("AAPL"); + + assertEquals(dummyPrice, result); + verify(cacheRepo, atLeastOnce()).find("AAPL"); + verify(dbRepo, times(1)).findLatest("AAPL"); + verify(cacheRepo, times(1)).saveIfAbsent("AAPL", dummyPrice); + verify(cacheRepo, never()).save(anyString(), any()); + } + + @Test + @DisplayName("캐시와 DB 모두에 가격 정보가 없으면 PriceNotFoundException을 던진다") + void shouldThrowExceptionIfNotFoundAnywhere() { + when(cacheRepo.find("AAPL")).thenReturn(Optional.empty()); + when(dbRepo.findLatest("AAPL")).thenReturn(Optional.empty()); + + PriceNotFoundException ex = assertThrows( + PriceNotFoundException.class, + () -> priceService.getCurrentPrice("AAPL") + ); + + assertEquals(PriceErrorCode.PRICE_NOT_FOUND, ex.getErrorCode()); + verify(cacheRepo, atLeastOnce()).find("AAPL"); + verify(dbRepo, times(1)).findLatest("AAPL"); + verify(cacheRepo, never()).saveIfAbsent(anyString(), any()); + verify(cacheRepo, never()).save(anyString(), any()); + } + + @Test + @DisplayName("캐시 적중 시 DB 접근 없이 즉시 반환하고 추가 쓰기 없음") + void cache_hit_returns_immediately_without_db_or_extra_writes() { + FakeRedisRepository redis = new FakeRedisRepository(); + FakeTimeSeriesRepository db = new FakeTimeSeriesRepository("AAPL", 111.0, 0); + redis.save("AAPL", new Price("AAPL", 111L, 0, 0.0, 0, ZonedDateTime.now())); + + var sut = new PriceQueryServiceImpl(redis, db, new SingleFlightExecutor<>(), executor, noBackoff, cacheProps); + + Price p = sut.getCurrentPrice("AAPL"); + + assertEquals(111.0, p.price(), 0.0); + assertEquals(1, redis.writeCount()); + assertEquals(0, db.reads.get()); + } + + @Test + @DisplayName("동시 캐시 미스 상황에서도 단일 DB 접근 및 단일 캐시 쓰기로 수렴") + void concurrent_cache_miss_collapses_to_single_write_via_single_flight_and_nx() throws Exception { + FakeRedisRepository redis = new FakeRedisRepository(); + FakeTimeSeriesRepository db = new FakeTimeSeriesRepository("AAPL", 123.45, 80); + var singleFlight = new SingleFlightExecutor(); + ExecutorService pool = Executors.newFixedThreadPool(64); + + var sut = new PriceQueryServiceImpl(redis, db, singleFlight, pool, noBackoff, cacheProps); + + int N = 200; + CountDownLatch ready = new CountDownLatch(N); + CountDownLatch start = new CountDownLatch(1); + CountDownLatch done = new CountDownLatch(N); + + for (int i = 0; i < N; i++) { + new Thread(() -> { + ready.countDown(); + try { + start.await(); + } catch (InterruptedException ignored) { + } + sut.getCurrentPrice("AAPL"); + done.countDown(); + }).start(); + } + + ready.await(); + start.countDown(); + done.await(); + + assertEquals(1, redis.writeCount()); + assertTrue(db.reads.get() <= 3); + pool.shutdownNow(); + } + + @Test + @DisplayName("백오프 중 다른 스레드가 캐시 채우면 재사용하고 쓰기 생략") + void when_peer_fills_cache_during_backoff_service_reuses_cache_and_skips_write() { + FakeRedisRepository redis = new FakeRedisRepository(); + FakeTimeSeriesRepository db = new FakeTimeSeriesRepository("AAPL", 200.0, 50); + + BackoffStrategy tinyBackoff = () -> { + try { + Thread.sleep(10); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + }; + + var service = new PriceQueryServiceImpl( + redis, db, new SingleFlightExecutor<>(), executor, tinyBackoff, new PriceCacheProperties(3) + ); + + new Thread(() -> { + try { + Thread.sleep(15); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + redis.saveIfAbsent("AAPL", new Price("AAPL", 200L, 200L, 0.0, 0, ZonedDateTime.now())); + }).start(); + + Price p = service.getCurrentPrice("AAPL"); + + assertEquals(200.0, p.price(), 0.0); + assertEquals(1, redis.writeCount()); + assertEquals(0, db.reads.get()); + } +} diff --git a/infra/build.gradle.kts b/infra/build.gradle.kts index 16b058c..9219957 100644 --- a/infra/build.gradle.kts +++ b/infra/build.gradle.kts @@ -18,9 +18,20 @@ repositories { } dependencies { + implementation(project(":domain:domain-price")) implementation("org.springframework.boot:spring-boot-starter") + implementation("org.springframework.boot:spring-boot-starter-web") + implementation("org.springframework:spring-jdbc") + runtimeOnly("org.postgresql:postgresql") + implementation("org.springframework.boot:spring-boot-starter-data-redis") + testImplementation("org.springframework.boot:spring-boot-starter-test") + testImplementation("com.h2database:h2") + testImplementation("org.springframework:spring-jdbc") testRuntimeOnly("org.junit.platform:junit-platform-launcher") + + compileOnly("org.projectlombok:lombok") + annotationProcessor("org.projectlombok:lombok") } tasks.withType { diff --git a/infra/src/main/java/com/polynomeer/infra/external/YahooFinancePriceClient.java b/infra/src/main/java/com/polynomeer/infra/external/YahooFinancePriceClient.java new file mode 100644 index 0000000..ce0892d --- /dev/null +++ b/infra/src/main/java/com/polynomeer/infra/external/YahooFinancePriceClient.java @@ -0,0 +1,27 @@ +package com.polynomeer.infra.external; + +import com.polynomeer.domain.price.model.Price; +import com.polynomeer.domain.price.repository.ExternalPriceClient; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Component; + +import java.time.ZonedDateTime; + +@Slf4j +@Component +public class YahooFinancePriceClient implements ExternalPriceClient { + + @Override + public Price fetchPrice(String tickerCode) { + log.info("Fetching price from YahooFinance API for {}", tickerCode); + + return new Price( + tickerCode, + 75000, + -100, + -0.13, + 12000000, + ZonedDateTime.now() + ); + } +} diff --git a/infra/src/main/java/com/polynomeer/infra/redis/RedisConfig.java b/infra/src/main/java/com/polynomeer/infra/redis/RedisConfig.java new file mode 100644 index 0000000..40b9134 --- /dev/null +++ b/infra/src/main/java/com/polynomeer/infra/redis/RedisConfig.java @@ -0,0 +1,34 @@ +package com.polynomeer.infra.redis; + +import com.fasterxml.jackson.databind.ObjectMapper; +import com.polynomeer.domain.price.model.Price; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.data.redis.connection.RedisConnectionFactory; +import org.springframework.data.redis.core.RedisTemplate; +import org.springframework.data.redis.serializer.Jackson2JsonRedisSerializer; +import org.springframework.data.redis.serializer.StringRedisSerializer; + +@Configuration +public class RedisConfig { + + @Bean + public RedisTemplate redisTemplate(RedisConnectionFactory connectionFactory) { + RedisTemplate template = new RedisTemplate<>(); + template.setConnectionFactory(connectionFactory); + template.setKeySerializer(new StringRedisSerializer()); + ObjectMapper objectMapper = new ObjectMapper().findAndRegisterModules(); + + Jackson2JsonRedisSerializer serializer = + new Jackson2JsonRedisSerializer<>(objectMapper, Price.class); + + template.setValueSerializer(serializer); + template.setHashKeySerializer(new StringRedisSerializer()); + template.setHashValueSerializer(serializer); + + template.afterPropertiesSet(); + return template; + } +} + + diff --git a/infra/src/main/java/com/polynomeer/infra/redis/RedisPriceRepository.java b/infra/src/main/java/com/polynomeer/infra/redis/RedisPriceRepository.java new file mode 100644 index 0000000..8bcacd8 --- /dev/null +++ b/infra/src/main/java/com/polynomeer/infra/redis/RedisPriceRepository.java @@ -0,0 +1,85 @@ +package com.polynomeer.infra.redis; + +import com.polynomeer.domain.price.model.Price; +import com.polynomeer.domain.price.repository.CachePriceRepository; +import jakarta.annotation.PostConstruct; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.data.redis.core.RedisTemplate; +import org.springframework.stereotype.Repository; + +import java.time.Duration; +import java.time.ZonedDateTime; +import java.util.Optional; + +@Slf4j +@Repository +@RequiredArgsConstructor +public class RedisPriceRepository implements CachePriceRepository { + + private static final Duration TTL = Duration.ofMinutes(1); + private final RedisTemplate redisTemplate; + + private String key(String tickerCode) { + return "price::" + tickerCode; + } + + @PostConstruct + public void saveTestData() { + Price price = new Price( + "BTC", + 50000L, + 100L, + 0.02, + 123456789L, + ZonedDateTime.now() + ); + redisTemplate.opsForValue().set("test:price:BTC", price); + } + + @Override + public Optional find(String tickerCode) { + String redisKey = key(tickerCode); + try { + Price cached = redisTemplate.opsForValue().get(redisKey); + if (cached != null) { + log.debug("[Redis] Cache hit: {}", redisKey); + } else { + log.debug("[Redis] Cache miss: {}", redisKey); + } + return Optional.ofNullable(cached); + } catch (Exception e) { + log.warn("[Redis] Cache read error for {}: {}", redisKey, e.getMessage()); + return Optional.empty(); + } + } + + @Override + public void save(String tickerCode, Price price) { + String redisKey = key(tickerCode); + try { + redisTemplate.opsForValue().set(redisKey, price, TTL); + log.debug("[Redis] Cache set: {} (TTL {}s)", redisKey, TTL.getSeconds()); + } catch (Exception e) { + log.warn("[Redis] Cache write error for {}: {}", redisKey, e.getMessage()); + } + } + + @Override + public boolean saveIfAbsent(String tickerCode, Price price) { + String redisKey = key(tickerCode); + try { + Boolean result = redisTemplate.opsForValue() + .setIfAbsent(redisKey, price, TTL); + boolean success = Boolean.TRUE.equals(result); + log.debug("[Redis] Cache {}: {} (TTL {}s)", + success ? "setIfAbsent" : "skip-existing", + redisKey, TTL.getSeconds()); + return success; + } catch (Exception e) { + log.warn("[Redis] Cache writeIfAbsent error for {}: {}", redisKey, e.getMessage()); + return false; + } + } + +} diff --git a/infra/src/main/java/com/polynomeer/infra/timescaledb/TimeSeriesChartRepositoryImpl.java b/infra/src/main/java/com/polynomeer/infra/timescaledb/TimeSeriesChartRepositoryImpl.java new file mode 100644 index 0000000..fbfa731 --- /dev/null +++ b/infra/src/main/java/com/polynomeer/infra/timescaledb/TimeSeriesChartRepositoryImpl.java @@ -0,0 +1,74 @@ +package com.polynomeer.infra.timescaledb; + +import com.polynomeer.domain.price.model.ChartPoint; +import com.polynomeer.domain.price.repository.TimeSeriesChartRepository; +import lombok.RequiredArgsConstructor; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Repository; + +import java.sql.PreparedStatement; +import java.time.OffsetDateTime; +import java.time.ZonedDateTime; +import java.util.List; + +@Repository +@RequiredArgsConstructor +public class TimeSeriesChartRepositoryImpl implements TimeSeriesChartRepository { + + private final JdbcTemplate jdbcTemplate; + + @Override + public List findChart(String tickerCode, String interval, ZonedDateTime from, ZonedDateTime to) { + String sql = """ + SELECT + time_bucket(?::interval, "timestamp") AS bucket, + first(price, "timestamp") AS open, + max(price) AS high, + min(price) AS low, + last(price, "timestamp") AS close, + sum(volume) AS volume + FROM price_history + WHERE ticker_code = ? + AND "timestamp" BETWEEN ? AND ? + GROUP BY bucket + ORDER BY bucket + """; + + String pgInterval = toPgInterval(interval); + OffsetDateTime fromTs = from.toOffsetDateTime(); + OffsetDateTime toTs = to.toOffsetDateTime(); + + return jdbcTemplate.query(con -> { + PreparedStatement ps = con.prepareStatement(sql); + ps.setString(1, pgInterval); + ps.setString(2, tickerCode); + ps.setObject(3, fromTs); + ps.setObject(4, toTs); + return ps; + }, (rs, i) -> new ChartPoint( + rs.getObject("bucket", OffsetDateTime.class).toZonedDateTime(), + rs.getLong("open"), + rs.getLong("high"), + rs.getLong("low"), + rs.getLong("close"), + rs.getLong("volume") + )); + } + + private String toPgInterval(String s) { + if (s == null || s.isEmpty()) throw new IllegalArgumentException("interval required"); + s = s.trim().toLowerCase(); + if (s.endsWith("m")) { + String n = s.substring(0, s.length() - 1); + return n + (n.equals("1") ? " minute" : " minutes"); + } else if (s.endsWith("h")) { + String n = s.substring(0, s.length() - 1); + return n + (n.equals("1") ? " hour" : " hours"); + } else if (s.endsWith("d")) { + String n = s.substring(0, s.length() - 1); + return n + (n.equals("1") ? " day" : " days"); + } else { + return s; + } + } +} diff --git a/infra/src/main/java/com/polynomeer/infra/timescaledb/TimeSeriesPriceRepositoryImpl.java b/infra/src/main/java/com/polynomeer/infra/timescaledb/TimeSeriesPriceRepositoryImpl.java new file mode 100644 index 0000000..431a077 --- /dev/null +++ b/infra/src/main/java/com/polynomeer/infra/timescaledb/TimeSeriesPriceRepositoryImpl.java @@ -0,0 +1,52 @@ +package com.polynomeer.infra.timescaledb; + +import com.polynomeer.domain.price.model.Price; +import com.polynomeer.domain.price.repository.TimeSeriesPriceRepository; +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.stereotype.Repository; + +import java.time.ZoneId; +import java.util.Optional; + +@Slf4j +@Repository +@RequiredArgsConstructor +public class TimeSeriesPriceRepositoryImpl implements TimeSeriesPriceRepository { + + private final JdbcTemplate jdbcTemplate; + + @Override + public Optional findLatest(String tickerCode) { + String sql = """ + SELECT price, volume, timestamp + FROM price_history + WHERE ticker_code = ? + ORDER BY timestamp DESC + LIMIT 1 + """; + + log.debug("[TimescaleDB] Executing SQL to find latest price for tickerCode={}", tickerCode); + + try { + Price price = jdbcTemplate.queryForObject(sql, new Object[]{tickerCode}, (rs, rowNum) -> { + Price p = new Price( + tickerCode, + rs.getLong("price"), + 0L, // change (not available) + 0.0, // changeRate (not available) + rs.getLong("volume"), + rs.getTimestamp("timestamp").toInstant().atZone(ZoneId.of("Asia/Seoul")) + ); + log.debug("[TimescaleDB] Query result mapped: {}", p); + return p; + }); + + return Optional.ofNullable(price); + } catch (Exception e) { + log.warn("[TimescaleDB] Failed to fetch latest price for {}: {}", tickerCode, e.getMessage()); + return Optional.empty(); + } + } +} diff --git a/infra/src/test/java/com/polynomeer/infra/redis/RedisPriceRepositoryTest.java b/infra/src/test/java/com/polynomeer/infra/redis/RedisPriceRepositoryTest.java new file mode 100644 index 0000000..9688cf1 --- /dev/null +++ b/infra/src/test/java/com/polynomeer/infra/redis/RedisPriceRepositoryTest.java @@ -0,0 +1,69 @@ +package com.polynomeer.infra.redis; + +import com.polynomeer.domain.price.model.Price; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.springframework.data.redis.core.RedisTemplate; +import org.springframework.data.redis.core.ValueOperations; + +import java.time.Duration; +import java.time.ZonedDateTime; +import java.util.Optional; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.*; + +class RedisPriceRepositoryTest { + + private RedisTemplate redisTemplate; + private RedisPriceRepository redisRepo; + + @BeforeEach + void setUp() { + redisTemplate = mock(RedisTemplate.class); + redisRepo = new RedisPriceRepository(redisTemplate); + } + + @Test + @DisplayName("정상적으로 캐시에 저장된다") + void shouldSavePriceToCache() { + var valueOps = mock(ValueOperations.class); + when(redisTemplate.opsForValue()).thenReturn(valueOps); + + Price price = dummyPrice(); + redisRepo.save("AAPL", price); + + verify(valueOps).set(eq("price::AAPL"), eq(price), any(Duration.class)); + } + + @Test + @DisplayName("캐시에서 Price 조회 성공 시 Optional.of 반환") + void shouldReturnOptionalOfPriceIfFound() { + var valueOps = mock(ValueOperations.class); + when(redisTemplate.opsForValue()).thenReturn(valueOps); + when(valueOps.get("price::AAPL")).thenReturn(dummyPrice()); + + Optional result = redisRepo.find("AAPL"); + + assertTrue(result.isPresent()); + assertEquals("AAPL", result.get().tickerCode()); + } + + @Test + @DisplayName("캐시에서 조회 실패 시 Optional.empty 반환") + void shouldReturnEmptyIfCacheMiss() { + var valueOps = mock(ValueOperations.class); + when(redisTemplate.opsForValue()).thenReturn(valueOps); + when(valueOps.get("price::AAPL")).thenReturn(null); + + Optional result = redisRepo.find("AAPL"); + + assertTrue(result.isEmpty()); + } + + private Price dummyPrice() { + return new Price("AAPL", 10000L, 100L, 1.0, 1000000L, ZonedDateTime.now()); + } +} diff --git a/infra/src/test/java/com/polynomeer/infra/timescaledb/TimeSeriesPriceRepositoryImplTest.java b/infra/src/test/java/com/polynomeer/infra/timescaledb/TimeSeriesPriceRepositoryImplTest.java new file mode 100644 index 0000000..3df349e --- /dev/null +++ b/infra/src/test/java/com/polynomeer/infra/timescaledb/TimeSeriesPriceRepositoryImplTest.java @@ -0,0 +1,64 @@ +package com.polynomeer.infra.timescaledb; + +import com.polynomeer.domain.price.model.Price; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.jdbc.datasource.DriverManagerDataSource; + +import java.util.Optional; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +class TimeSeriesPriceRepositoryImplTest { + + private TimeSeriesPriceRepositoryImpl timeSeriesRepo; + + @BeforeEach + void setUp() { + var dataSource = new DriverManagerDataSource(); + dataSource.setDriverClassName("org.h2.Driver"); + dataSource.setUrl("jdbc:h2:mem:testdb;MODE=PostgreSQL;DB_CLOSE_DELAY=-1"); + dataSource.setUsername("sa"); + dataSource.setPassword(""); + + JdbcTemplate jdbcTemplate = new JdbcTemplate(dataSource); + timeSeriesRepo = new TimeSeriesPriceRepositoryImpl(jdbcTemplate); + + jdbcTemplate.execute("DROP TABLE IF EXISTS price_history"); + jdbcTemplate.execute(""" + CREATE TABLE price_history ( + ticker_code VARCHAR(20), + timestamp TIMESTAMP, + price BIGINT, + volume BIGINT + ) + """); + + // 해외 주식 예시 데이터 삽입 (AAPL) + jdbcTemplate.execute(""" + INSERT INTO price_history (ticker_code, timestamp, price, volume) VALUES + ('AAPL', '2024-08-01T09:00:00', 19500, 500000) + """); + } + + @Test + @DisplayName("해외 종목 AAPL의 최근 시세를 반환한다") + void shouldReturnLatestPriceForAAPL() { + Optional result = timeSeriesRepo.findLatest("AAPL"); + + assertTrue(result.isPresent()); + assertEquals("AAPL", result.get().tickerCode()); + assertEquals(19500, result.get().price()); + } + + @Test + @DisplayName("존재하지 않는 해외 종목에 대해 Optional.empty 반환") + void shouldReturnEmptyForInvalidTicker() { + Optional result = timeSeriesRepo.findLatest("TSLA"); + assertTrue(result.isEmpty()); + } +} + diff --git a/shared/build.gradle.kts b/shared/build.gradle.kts index 7e3a7e3..da4008d 100644 --- a/shared/build.gradle.kts +++ b/shared/build.gradle.kts @@ -20,6 +20,7 @@ repositories { dependencies { implementation("org.springframework.boot:spring-boot-starter") implementation("org.springframework.boot:spring-boot-starter-web") + implementation("org.springframework.boot:spring-boot-starter-data-redis") implementation("org.springdoc:springdoc-openapi-starter-webmvc-api:2.8.9") testImplementation("org.springframework.boot:spring-boot-starter-test") testRuntimeOnly("org.junit.platform:junit-platform-launcher") diff --git a/shared/src/main/java/com/polynomeer/shared/common/error/PriceErrorCode.java b/shared/src/main/java/com/polynomeer/shared/common/error/PriceErrorCode.java new file mode 100644 index 0000000..3143b05 --- /dev/null +++ b/shared/src/main/java/com/polynomeer/shared/common/error/PriceErrorCode.java @@ -0,0 +1,15 @@ +package com.polynomeer.shared.common.error; + +import lombok.AllArgsConstructor; +import lombok.Getter; +import org.springframework.http.HttpStatus; + +@AllArgsConstructor +@Getter +public enum PriceErrorCode implements ErrorCode { + PRICE_NOT_FOUND("PRICE-001", "해당 종목의 가격을 찾을 수 없습니다.", HttpStatus.NOT_FOUND); + + private final String code; + private final String message; + private final HttpStatus httpStatus; +} diff --git a/shared/src/main/java/com/polynomeer/shared/common/error/PriceNotFoundException.java b/shared/src/main/java/com/polynomeer/shared/common/error/PriceNotFoundException.java new file mode 100644 index 0000000..3d9e1df --- /dev/null +++ b/shared/src/main/java/com/polynomeer/shared/common/error/PriceNotFoundException.java @@ -0,0 +1,9 @@ +package com.polynomeer.shared.common.error; + +public class PriceNotFoundException extends BaseException { + public PriceNotFoundException(ErrorCode errorCode) { + super(PriceErrorCode.PRICE_NOT_FOUND); + } + + +} diff --git a/shared/src/main/java/com/polynomeer/shared/common/error/TickerErrorCode.java b/shared/src/main/java/com/polynomeer/shared/common/error/TickerErrorCode.java index 0efc96f..ed529f4 100644 --- a/shared/src/main/java/com/polynomeer/shared/common/error/TickerErrorCode.java +++ b/shared/src/main/java/com/polynomeer/shared/common/error/TickerErrorCode.java @@ -6,7 +6,8 @@ @AllArgsConstructor public enum TickerErrorCode implements ErrorCode { - TICKER_NOT_FOUND("TICKER-001", "해당 종목을 찾을 수 없습니다.", HttpStatus.NOT_FOUND); + TICKER_NOT_FOUND("TICKER-001", "해당 종목을 찾을 수 없습니다.", HttpStatus.NOT_FOUND), + TICKER_INVALID("TICKER-002", "요청된 Ticker가 유효하지 않습니다.", HttpStatus.BAD_REQUEST); @Getter private final String code; diff --git a/shared/src/main/java/com/polynomeer/shared/common/error/TickerValidationException.java b/shared/src/main/java/com/polynomeer/shared/common/error/TickerValidationException.java new file mode 100644 index 0000000..806aba7 --- /dev/null +++ b/shared/src/main/java/com/polynomeer/shared/common/error/TickerValidationException.java @@ -0,0 +1,7 @@ +package com.polynomeer.shared.common.error; + +public class TickerValidationException extends BaseException { + public TickerValidationException(ErrorCode errorCode) { + super(TickerErrorCode.TICKER_INVALID); + } +}