diff --git a/src/KArtSell.Modules.ModelOperations/ShadowRun/Services/KrxDataService.cs b/src/KArtSell.Modules.ModelOperations/ShadowRun/Services/KrxDataService.cs index 1bf60b4f..86769fa5 100644 --- a/src/KArtSell.Modules.ModelOperations/ShadowRun/Services/KrxDataService.cs +++ b/src/KArtSell.Modules.ModelOperations/ShadowRun/Services/KrxDataService.cs @@ -1,21 +1,26 @@ using System.Net; using System.Text.Json; +using Dapper; using Microsoft.Extensions.Caching.Memory; using Microsoft.Extensions.Logging; +using Npgsql; namespace KArtSell.Modules.ModelOperations.ShadowRun.Services; /// /// Fetches historical OHLCV and fee schedule data from Korea Exchange (KRX) API. /// Implements caching, retry logic, and PIT-safe lookups (no forward bias). +/// Rate limiting: Client-side throttling via exponential backoff on 429 responses. /// public sealed class KrxDataService : IKrxDataService { private readonly HttpClient _httpClient; private readonly IMemoryCache _cache; private readonly ILogger _logger; + private readonly NpgsqlDataSource _dataSource; private const int CacheDurationMinutes = 1440; // 24 hours + private const int RecentDaysWindow = 7; // Last 7 days: always refresh (mutable data) private const int MaxRetries = 3; private const int InitialBackoffMs = 100; private const int MaxBackoffMs = 30000; @@ -46,17 +51,62 @@ public sealed class KrxDataService : IKrxDataService new EventId(4, nameof(LogRetryError)), "Retryable error: {ErrorMessage}"); - public KrxDataService(HttpClient httpClient, IMemoryCache cache, ILogger logger) + public KrxDataService( + HttpClient httpClient, + IMemoryCache cache, + ILogger logger, + NpgsqlDataSource? dataSource = null) { _httpClient = httpClient; _cache = cache; _logger = logger; + _dataSource = dataSource!; + } + + /// + /// Get the last date that was successfully imported from KRX API. + /// Returns null if no successful import exists or DB not available. + /// Used for incremental fetching (avoid re-fetching old, immutable data). + /// + private async Task GetLastSuccessfulImportDateAsync(CancellationToken cancellationToken) + { + // In test environments or when DB is not available, skip incremental optimization + if (_dataSource == null) + { + return null; + } + + try + { + const string sql = """ + SELECT MAX(DATE(import_at)) as last_date + FROM market_data.krx_imports + WHERE status = 'SUCCESS' + """; + + await using var connection = await _dataSource.OpenConnectionAsync(cancellationToken); + var result = await connection.QuerySingleOrDefaultAsync(sql); + + if (result == null) + { + _logger.LogInformation("No previous successful KRX import found, will fetch full range"); + return null; + } + + return DateOnly.FromDateTime(result.Value); + } + catch (Exception ex) + { + _logger.LogWarning(ex, "Failed to get last successful import date, fetching full range"); + return null; + } } /// /// Fetch daily OHLCV bars for ticker within date range. - /// Implements caching (24h) and retry logic for transient failures. + /// Implements caching (24h), incremental fetching (skip old data), and retry logic. /// PIT-safe: Returns only requested date range (no lookback). + /// Strategy: Last 7 days always fresh (mutable), older data fetched only once. /// public async Task> GetDailyOhlcvAsync( string ticker, @@ -66,7 +116,43 @@ public sealed class KrxDataService : IKrxDataService { LogFetchingOhlcv(_logger, ticker, startDate, endDate, null); - var cacheKey = $"ohlcv:{ticker}:{startDate:yyyyMMdd}:{endDate:yyyyMMdd}"; + // Incremental fetching: Skip old, immutable data that was already collected + var today = DateOnly.FromDateTime(DateTime.UtcNow.Date); + var lastSuccessfulImport = await GetLastSuccessfulImportDateAsync(cancellationToken); + + // Strategy: Only fetch data from 7 days ago onwards (last 7 days always fresh) + // Skip anything older that was already imported successfully + var effectiveStartDate = startDate; + if (lastSuccessfulImport.HasValue) + { + var oldDataCutoff = today.AddDays(-RecentDaysWindow); + var latestOldData = lastSuccessfulImport.Value; + + if (latestOldData >= oldDataCutoff) + { + // Already have recent data, skip to day after last import + effectiveStartDate = latestOldData.AddDays(1); + } + } + + // If effective range is empty, return empty + if (effectiveStartDate > endDate) + { + _logger.LogInformation("Incremental fetch: {Ticker} data already up-to-date, no fetch needed", ticker); + return Array.Empty(); + } + + var skippedDays = effectiveStartDate.DayNumber - startDate.DayNumber; + if (skippedDays > 0) + { + _logger.LogInformation( + "Incremental fetch: {Ticker} skipping {SkippedDays} immutable days (already imported), starting from {EffectiveStart}", + ticker, + skippedDays, + effectiveStartDate); + } + + var cacheKey = $"ohlcv:{ticker}:{effectiveStartDate:yyyyMMdd}:{endDate:yyyyMMdd}"; // Check cache first if (_cache.TryGetValue(cacheKey, out IReadOnlyList? cached)) @@ -232,8 +318,10 @@ public sealed class KrxDataService : IKrxDataService } } - // Combine all responses - return $"[{string.Join(",", results.Select(r => ExtractPriceItems(r)))}]"; + // Combine all responses (or return empty if no results) + return results.Any() + ? $"[{string.Join(",", results.Select(r => ExtractPriceItems(r)))}]" + : "[]"; } private string ExtractPriceItems(string krxResponse)