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;
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;
}
using var scope = scopeFactory.CreateScope();
var db = scope.ServiceProvider.GetRequiredService();
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;
}
}