using Application.Abstractions.Data; using Application.Helpers; using Domain.Entities.Stocks; using Domain.Entities.Stocks.ValueObject; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using SharedKernel; namespace Infrastructure.StockData; /// /// 거시·환율 수집 (한국수출입은행 OpenAPI) — AP01 현재환율 + AP02 대출금리 + AP03 국제금리. 기본 11:30 KST 실행 /// (koreaexim 은 영업일 ~11시 이후 반영). 공용 KrxBackfill 로 [today-MacroBackfillDays, 직전 영업일] 을 최신일부터 과거로 훑으며 /// 미적재일만 채운다 — quota 보호를 위해 1회 실행당 BackfillMaxPerRun(기본 60)일까지만 fetch. 각 날짜는 세 data 종류를 /// (키가 설정된 것만) 모두 수집해 ExchangeRate / InterestRate upsert (UQ = (CurrencyUnit,TradeDate) / (RateType,ItemName,TradeDate)). /// 빈 배열([]) — 주말/공휴일·오전 미반영 — 은 "데이터 없음" 으로 취급(그 날은 미적재 유지, 다음 실행 때 재시도). 키 전부 미설정 시 skip. /// existsForDate 판정은 세 종류 중 하나라도 그 날짜에 적재됐는지로 본다 (부분 성공 후 재실행 시 남은 종류를 다시 채울 수 있게 upsert 로 안전). /// internal sealed class MacroDataSyncService( IServiceScopeFactory scopeFactory, IHttpClientFactory httpClientFactory, IOptions settings, ILogger logger ) : DailyScheduledService(logger) { protected override string JobName => "MacroDataSync"; protected override TimeOnly TargetTime => ParseTime(settings.Value.KoreaExim.MacroSyncTime, new TimeOnly(11, 30)); protected override int MaxRetryCount => 2; protected override TimeSpan RetryDelay => TimeSpan.FromHours(1); protected override async Task RunOnceAsync(DateOnly todayKst, CancellationToken ct) { var cfg = settings.Value.KoreaExim; using var scope = scopeFactory.CreateScope(); var collectorSettings = scope.ServiceProvider.GetRequiredService(); if (!await collectorSettings.IsEnabledAsync(CollectorFlag.KoreaEximMacro, ct)) { return true; } cfg = cfg with { ExchangeKey = await collectorSettings.GetKeyAsync(CollectorKey.KoreaEximExchange, ct) ?? cfg.ExchangeKey, LoanRateKey = await collectorSettings.GetKeyAsync(CollectorKey.KoreaEximLoanRate, ct) ?? cfg.LoanRateKey, IntlRateKey = await collectorSettings.GetKeyAsync(CollectorKey.KoreaEximIntlRate, ct) ?? cfg.IntlRateKey }; var db = scope.ServiceProvider.GetRequiredService(); var hasExchange = !string.IsNullOrWhiteSpace(cfg.ExchangeKey); var hasLoan = !string.IsNullOrWhiteSpace(cfg.LoanRateKey); var hasIntl = !string.IsNullOrWhiteSpace(cfg.IntlRateKey); if (!hasExchange && !hasLoan && !hasIntl) { Logger.LogWarning("[{Job}] KoreaExim 키(ExchangeKey/LoanRateKey/IntlRateKey) 전부 미설정 — 수집 skip", JobName); return true; } var client = httpClientFactory.CreateClient(KoreaEximHttp.ClientName); var endDate = await MarketCalendar.GetPreviousBusinessDayAsync(db, todayKst.AddDays(1), ct); var days = cfg.MacroBackfillDays > 0 ? cfg.MacroBackfillDays : 30; var startDate = todayKst.AddDays(-days); var holidays = (await db.MarketHoliday.AsNoTracking() .Where(c => c.Date >= startDate && c.Date <= endDate) .Select(c => c.Date) .ToListAsync(ct)).ToHashSet(); var maxPerRun = cfg.BackfillMaxPerRun > 0 ? cfg.BackfillMaxPerRun : 60; var baseUrl = cfg.BaseUrl.TrimEnd('/'); var fetched = await KrxBackfill.RunAsync( existsForDate: (day, token) => ExistsForDateAsync(db, day, hasExchange, token), fetchAndUpsertForDate: (day, token) => FetchAndUpsertAsync(db, client, baseUrl, cfg, hasExchange, hasLoan, hasIntl, day, token), startDate: startDate, endDate: endDate, holidays: holidays, maxPerRun: maxPerRun, delayMs: 300, ct: ct); Logger.LogInformation("[{Job}] 완료 — 창=[{Start}~{End}], 이번 실행 fetch={Fetched}일 (maxPerRun={Max})", JobName, startDate, endDate, fetched, maxPerRun); return true; } /// 해당 날짜 적재 여부 — 환율 키가 있으면 ExchangeRate, 없으면 InterestRate 로 판정 private static async Task ExistsForDateAsync(IAppDbContext db, DateOnly day, bool hasExchange, CancellationToken ct) { if (hasExchange) { return await db.ExchangeRate.AsNoTracking().AnyAsync(c => c.TradeDate == day, ct); } return await db.InterestRate.AsNoTracking().AnyAsync(c => c.TradeDate == day, ct); } private async Task FetchAndUpsertAsync( IAppDbContext db, HttpClient client, string baseUrl, AppSettings.KoreaEximSection cfg, bool hasExchange, bool hasLoan, bool hasIntl, DateOnly day, CancellationToken ct ) { var touched = 0; if (hasExchange) { touched += await FetchExchangeAsync(db, client, baseUrl, cfg.ExchangeKey, day, ct); } if (hasLoan) { touched += await FetchRateAsync(db, client, baseUrl, cfg.LoanRateKey, "AP02", RateType.Loan, day, ct); } if (hasIntl) { touched += await FetchRateAsync(db, client, baseUrl, cfg.IntlRateKey, "AP03", RateType.International, day, ct); } if (touched > 0) { await db.SaveChangesAsync(ct); } else { Logger.LogInformation("[{Job}] searchdate={Day} 거시 미반영 (0건)", JobName, day); } } /// AP01 환율 수집 + upsert. 반환값 = 이번 호출로 add/update 한 행 수. private async Task FetchExchangeAsync(IAppDbContext db, HttpClient client, string baseUrl, string authKey, DateOnly day, CancellationToken ct) { var url = $"{baseUrl}?authkey={authKey}&searchdate={day:yyyyMMdd}&data=AP01"; var json = await KoreaEximHttp.GetStringWithRetryAsync(client, url, Logger, ct); var parsed = KoreaEximExchangeParser.Parse(json); Logger.LogInformation("[{Job}] AP01 환율 searchdate={Day} rows={Rows}", JobName, day, parsed.Count); if (parsed.Count == 0) { return 0; } var existing = await db.ExchangeRate.Where(c => c.TradeDate == day).ToListAsync(ct); var existingByUnit = existing.ToDictionary(c => c.CurrencyUnit); var count = 0; foreach (var row in parsed) { if (existingByUnit.TryGetValue(row.CurrencyUnit, out var rate)) { rate.Update(row.CurrencyName, row.DealBasRate, row.Ttb, row.Tts); } else { var created = ExchangeRate.Create(row.CurrencyUnit, row.CurrencyName, day, row.DealBasRate, row.Ttb, row.Tts); await db.ExchangeRate.AddAsync(created, ct); existingByUnit[row.CurrencyUnit] = created; } count++; } return count; } /// AP02/AP03 금리 수집 + upsert. 반환값 = 이번 호출로 add/update 한 행 수. private async Task FetchRateAsync(IAppDbContext db, HttpClient client, string baseUrl, string authKey, string data, RateType rateType, DateOnly day, CancellationToken ct) { var url = $"{baseUrl}?authkey={authKey}&searchdate={day:yyyyMMdd}&data={data}"; var json = await KoreaEximHttp.GetStringWithRetryAsync(client, url, Logger, ct); var parsed = KoreaEximRateParser.Parse(json); Logger.LogInformation("[{Job}] {Data} {Type} searchdate={Day} rows={Rows}", JobName, data, rateType, day, parsed.Count); if (parsed.Count == 0) { return 0; } var existing = await db.InterestRate.Where(c => c.RateType == rateType && c.TradeDate == day).ToListAsync(ct); var existingByName = existing.ToDictionary(c => c.ItemName); var count = 0; foreach (var row in parsed) { if (existingByName.TryGetValue(row.ItemName, out var rate)) { rate.Update(row.Rate); } else { var created = InterestRate.Create(rateType, row.ItemName, day, row.Rate); await db.InterestRate.AddAsync(created, ct); existingByName[row.ItemName] = created; } count++; } return count; } }