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; /// /// 종목 마스터 동기화 — 금융위 KRX 상장종목정보 API 를 일 1회(기본 07:30 KST) 전량 수집하여 Stock upsert. /// 신규 상장 = insert, 명칭/시장 변경 = update, 스냅샷에서 사라진 종목 = 상폐 soft-off. /// 단일 엔드포인트 페이징 수집이라 조기 종료/페이지 상한 시 부분 스냅샷이 될 수 있으므로, /// 상폐 sweep 은 시장별 가드(StockMasterDelist — 0건 또는 활성 대비 80% 미만이면 보류+경고)를 통과한 시장에서만 수행. /// ServiceKey 미설정 시 로그만 남기고 skip. /// internal sealed class StockMasterSyncService( IServiceScopeFactory scopeFactory, IHttpClientFactory httpClientFactory, IOptions settings, ILogger logger ) : DailyScheduledService(logger) { private const string ServicePath = "/1160100/service/GetKrxListedInfoService/getItemInfo"; private const int MaxPages = 50; private const int MaxBaseDateLookback = 7; protected override string JobName => "StockMasterSync"; protected override TimeOnly TargetTime => ParseTime(settings.Value.StockData.MasterSyncTime, new TimeOnly(7, 30)); protected override async Task RunOnceAsync(DateOnly todayKst, CancellationToken ct) { var cfg = settings.Value.StockData.DataGoKr; using var scope = scopeFactory.CreateScope(); var collectorSettings = scope.ServiceProvider.GetRequiredService(); if (!await collectorSettings.IsEnabledAsync(CollectorFlag.StockDataMaster, ct)) { return true; } cfg = cfg with { ServiceKey = await collectorSettings.GetKeyAsync(CollectorKey.DataGoKr, ct) ?? cfg.ServiceKey }; if (string.IsNullOrWhiteSpace(cfg.ServiceKey)) { Logger.LogWarning("[{Job}] StockData:DataGoKr:ServiceKey 미설정 — 수집 skip", JobName); return true; } var db = scope.ServiceProvider.GetRequiredService(); var client = httpClientFactory.CreateClient(DataGoKrHttp.ClientName); // 상장종목정보는 basDt 단위 스냅샷 — 직전 영업일부터 최대 7일 소급하며 데이터가 있는 기준일을 찾는다 var baseDate = await MarketCalendar.GetPreviousBusinessDayAsync(db, todayKst.AddDays(1), ct); List? snapshot = null; for (var back = 0; back < MaxBaseDateLookback; back++) { snapshot = await FetchSnapshotAsync(client, cfg, baseDate, ct); if (snapshot.Count > 0) { break; } baseDate = await MarketCalendar.GetPreviousBusinessDayAsync(db, baseDate, ct); } if (snapshot is null || snapshot.Count == 0) { Logger.LogError("[{Job}] 상장종목 스냅샷 없음 — 최근 {Days}영업일 조회 실패", JobName, MaxBaseDateLookback); return true; } // 동일 코드 중복 행은 마지막 행 우선. 상폐 sweep 가드용으로 시장별 건수도 함께 집계 (0건 시장 포함) var byCode = new Dictionary(); var snapshotCountByMarket = new Dictionary { [StockMarket.KOSPI] = 0, [StockMarket.KOSDAQ] = 0, [StockMarket.KONEX] = 0 }; foreach (var item in snapshot) { var market = MapMarket(item.MarketName); if (market is null) { continue; } if (!byCode.ContainsKey(item.Code)) { snapshotCountByMarket[market.Value]++; } byCode[item.Code] = item; } var stocks = await db.Stock.ToListAsync(ct); var stockByCode = stocks.ToDictionary(c => c.Code); var inserted = 0; var updated = 0; var delisted = 0; foreach (var (code, item) in byCode) { var market = MapMarket(item.MarketName)!.Value; if (stockByCode.TryGetValue(code, out var stock)) { var wasUpdated = stock.UpdatedAt; stock.UpdateMaster(item.Name, market, item.Isin); if (stock.UpdatedAt != wasUpdated) { updated++; } } else { await db.Stock.AddAsync(Stock.Create(code, item.Name, market, item.BaseDate, item.Isin), ct); inserted++; } } // 스냅샷에서 사라진 활성 종목 = 상폐 soft-off — 단, 시장별 가드(빈 응답/활성 대비 80% 미만)를 통과한 시장만 var activeCountByMarket = stocks.Where(c => c.IsActive).GroupBy(c => c.Market).ToDictionary(c => c.Key, c => c.Count()); var delistableMarkets = StockMasterDelist.GetDelistableMarkets(snapshotCountByMarket, activeCountByMarket).ToHashSet(); foreach (var (market, snapshotCount) in snapshotCountByMarket.Where(c => !delistableMarkets.Contains(c.Key))) { Logger.LogWarning("[{Job}] {Market} 상폐 sweep 보류 — snapshot={Snapshot}, active={Active} (빈/부분 응답 방어)", JobName, market, snapshotCount, activeCountByMarket.TryGetValue(market, out var active) ? active : 0); } foreach (var stock in stocks.Where(c => c.IsActive && delistableMarkets.Contains(c.Market) && !byCode.ContainsKey(c.Code))) { stock.MarkDelisted(baseDate); delisted++; } await db.SaveChangesAsync(ct); Logger.LogInformation("[{Job}] 완료 — basDt={BaseDate}, snapshot={Snapshot}, inserted={Inserted}, updated={Updated}, delisted={Delisted}", JobName, baseDate, byCode.Count, inserted, updated, delisted); return true; } private async Task> FetchSnapshotAsync(HttpClient client, AppSettings.StockDataSection.DataGoKrSection cfg, DateOnly baseDate, CancellationToken ct) { var all = new List(); var totalCount = int.MaxValue; for (var pageNo = 1; pageNo <= MaxPages && all.Count < totalCount; pageNo++) { var url = $"{cfg.BaseUrl.TrimEnd('/')}{ServicePath}?serviceKey={Uri.EscapeDataString(cfg.ServiceKey)}&resultType=json&numOfRows={cfg.PageSize}&pageNo={pageNo}&basDt={baseDate:yyyyMMdd}"; var json = await DataGoKrHttp.GetStringWithRetryAsync(client, url, Logger, ct); var (items, total) = DataGoKrStockParser.ParseListedItems(json); totalCount = total; if (items.Count == 0) { break; } all.AddRange(items); } return all; } private static StockMarket? MapMarket(string marketName) { if (marketName.Contains("KOSPI", StringComparison.OrdinalIgnoreCase)) { return StockMarket.KOSPI; } if (marketName.Contains("KOSDAQ", StringComparison.OrdinalIgnoreCase)) { return StockMarket.KOSDAQ; } if (marketName.Contains("KONEX", StringComparison.OrdinalIgnoreCase)) { return StockMarket.KONEX; } return null; } }