using Application.Abstractions.Data; using Application.Helpers; using Domain.Entities.Stocks; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using SharedKernel; namespace Infrastructure.StockData; /// /// 전자공시(DART) 공시 목록 수집 (OpenDART list.json) — 기본 08:30 KST 실행. 접수일 [today-BackfillDays, today] 범위를 /// page_count=100 으로 페이징(total_page 만큼 순회)하며 수집한다. KRX 배치처럼 날짜별 반복이 아니라 DART 는 날짜범위 페이징이므로 /// KrxBackfill 대신 자체 페이징 루프를 쓴다. 접수번호(rcept_no) 로 upsert 하되 이미 있으면 skip(idempotent) — 갱신 없이 신규만 삽입. /// • 스펙 클램프: corp_code 미지정 검색은 최근 3개월만 제공 — BackfillDays 가 90 초과여도 창을 3개월로 클램프(경고). /// • 초기 백필(창의 가장 오래된 주에 적재분 0건): bgn_de/end_de 를 주 단위로 분할 순회해 5만 행(MaxPages×100) 상한을 회피. /// • 증분 수집: 한 페이지 전체가 기존 건이면(1페이지 제외) 그보다 오래된 페이지도 기수집으로 보고 조기 종료 — 매일 전량 재페이징 방지. /// • 페이지 상한 도달 시 잔여 페이지 수를 경고 (정렬 desc 라 오래된 공시부터 유실). /// status "013"(데이터 없음) 은 정상 종료, "000" 외 다른 status 는 오류로 보고 후 중단(재시도 대상 — DescribeStatus 라벨 로깅). /// ApiKey 미설정 시 로그만 남기고 skip. list.json 이 stock_code 를 포함하므로 corpCode.xml 매핑 불필요. /// internal sealed class DisclosureSyncService( IServiceScopeFactory scopeFactory, IHttpClientFactory httpClientFactory, IOptions settings, ILogger logger ) : DailyScheduledService(logger) { private const int PageCount = 100; private const int MaxPages = 500; // 안전 상한 (total_page 무한루프 방지) protected override string JobName => "DisclosureSync"; protected override TimeOnly TargetTime => ParseTime(settings.Value.OpenDart.DisclosureSyncTime, new TimeOnly(8, 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.OpenDart; using var scope = scopeFactory.CreateScope(); var collectorSettings = scope.ServiceProvider.GetRequiredService(); if (!await collectorSettings.IsEnabledAsync(CollectorFlag.OpenDartDisclosure, ct)) { return true; } cfg = cfg with { ApiKey = await collectorSettings.GetKeyAsync(CollectorKey.OpenDart, ct) ?? cfg.ApiKey }; if (string.IsNullOrWhiteSpace(cfg.ApiKey)) { Logger.LogWarning("[{Job}] OpenDart:ApiKey 미설정 — 수집 skip", JobName); return true; } var db = scope.ServiceProvider.GetRequiredService(); var client = httpClientFactory.CreateClient(OpenDartHttp.ClientName); var days = cfg.BackfillDays > 0 ? cfg.BackfillDays : 30; var beginDate = todayKst.AddDays(-days); // corp_code 미지정 list.json 은 스펙상 최근 3개월만 제공 — 초과 설정 시 클램프 var (clampedBegin, clamped) = DartBackfill.ClampToSpecWindow(beginDate, todayKst); if (clamped) { Logger.LogWarning("[{Job}] BackfillDays={Days} 가 스펙 한도(최근 {Months}개월)를 초과 — 창 시작일을 {Begin} 으로 클램프", JobName, days, DartBackfill.SpecWindowMonths, clampedBegin); beginDate = clampedBegin; } var baseUrl = cfg.BaseUrl.TrimEnd('/'); // 초기 백필 판정 — 창의 가장 오래된 주(7일)에 적재된 공시가 하나도 없으면 미커버로 간주 (영업일마다 수백 건이 접수되므로 안전) var oldestWeekEnd = beginDate.AddDays(7); var isInitialBackfill = !await db.Disclosure.AsNoTracking().AnyAsync(c => c.RceptDt >= beginDate && c.RceptDt < oldestWeekEnd, ct); var inserted = 0; var skipped = 0; var ok = true; if (isInitialBackfill) { // 최초 적재: 창 전체가 5만 행(MaxPages×PageCount) 상한을 넘을 수 있어 주 단위로 분할 순회 (사업보고서 시즌 30일 창이 5만 초과 현실적) var chunks = DartBackfill.SplitWeekly(beginDate, todayKst); Logger.LogInformation("[{Job}] 초기 백필 — 범위=[{Begin}~{Today}] 를 주 단위 {Chunks}개 창으로 분할 수집", JobName, beginDate, todayKst, chunks.Count); foreach (var (chunkBegin, chunkEnd) in chunks) { var result = await SyncRangeAsync(db, client, baseUrl, cfg.ApiKey, chunkBegin, chunkEnd, allowEarlyStop: false, ct); inserted += result.Inserted; skipped += result.Skipped; if (!result.Ok) { ok = false; break; } await Task.Delay(300, ct); // quota 보호 (청크 간) } } else { // 증분 수집: 신규분이 최신 페이지에 몰려 있으므로 전체 기존 페이지를 만나면 조기 종료 var result = await SyncRangeAsync(db, client, baseUrl, cfg.ApiKey, beginDate, todayKst, allowEarlyStop: true, ct); inserted = result.Inserted; skipped = result.Skipped; ok = result.Ok; } Logger.LogInformation("[{Job}] 완료 — 범위=[{Begin}~{Today}], inserted={Inserted}, skipped={Skipped}, 초기백필={Initial}", JobName, beginDate, todayKst, inserted, skipped, isInitialBackfill); return ok; } /// /// [bgnDe, endDe] 접수일 범위를 desc 정렬로 페이징 수집. 신규만 삽입(rcept_no idempotent). /// allowEarlyStop=true 면 "페이지 전체 기존 + 1페이지 아님" 조건에서 조기 종료 (DartBackfill.ShouldStopEarly). /// API 오류 status 는 Ok=false 로 반환해 베이스 재시도로 연결한다. /// private async Task<(bool Ok, int Inserted, int Skipped)> SyncRangeAsync( IAppDbContext db, HttpClient client, string baseUrl, string apiKey, DateOnly bgnDe, DateOnly endDe, bool allowEarlyStop, CancellationToken ct) { var inserted = 0; var skipped = 0; var page = 1; var totalPage = 1; while (page <= totalPage && page <= MaxPages) { ct.ThrowIfCancellationRequested(); var url = $"{baseUrl}/api/list.json?crtfc_key={apiKey}&bgn_de={bgnDe:yyyyMMdd}&end_de={endDe:yyyyMMdd}&page_no={page}&page_count={PageCount}"; var json = await OpenDartHttp.GetStringWithRetryAsync(client, url, Logger, ct); var result = DartDisclosureParser.Parse(json); if (result.Status == "013") { Logger.LogInformation("[{Job}] 조회된 데이터 없음(013) — 종료. 범위=[{Begin}~{End}]", JobName, bgnDe, endDe); break; } if (result.Status != "000") { Logger.LogError("[{Job}] OpenDART 오류 status={Status}({Label}) message={Message} (page={Page}, 범위=[{Begin}~{End}])", JobName, result.Status, DartDisclosureParser.DescribeStatus(result.Status), result.Message, page, bgnDe, endDe); return (false, inserted, skipped); } totalPage = result.TotalPage > 0 ? result.TotalPage : 1; if (result.Rows.Count == 0) { break; } var rceptNos = result.Rows.Select(c => c.RceptNo).ToList(); var existing = (await db.Disclosure.AsNoTracking() .Where(c => rceptNos.Contains(c.RceptNo)) .Select(c => c.RceptNo) .ToListAsync(ct)).ToHashSet(); var newCount = 0; foreach (var row in result.Rows) { if (existing.Contains(row.RceptNo)) { skipped++; continue; } var created = Disclosure.Create(row.RceptNo, row.CorpCode, row.CorpName, row.StockCode, row.CorpCls, row.ReportNm, row.FlrNm, row.RceptDt, row.Rm); await db.Disclosure.AddAsync(created, ct); existing.Add(row.RceptNo); // 같은 실행 내 중복 접수번호 방어 inserted++; newCount++; } await db.SaveChangesAsync(ct); Logger.LogInformation("[{Job}] page {Page}/{TotalPage} 처리 — rows={Rows}, 신규={New} (누적 inserted={Inserted}, skipped={Skipped})", JobName, page, totalPage, result.Rows.Count, newCount, inserted, skipped); if (allowEarlyStop && DartBackfill.ShouldStopEarly(page, result.Rows.Count, newCount)) { Logger.LogInformation("[{Job}] 조기 종료 — page {Page} 전체가 기수집 건 (이후 페이지는 더 오래된 기수집분)", JobName, page); return (true, inserted, skipped); } page++; if (page <= totalPage) { await Task.Delay(300, ct); // quota 보호 } } if (totalPage > MaxPages) { Logger.LogWarning("[{Job}] 페이지 상한 도달 — totalPage={TotalPage}, 상한={MaxPages}, 미수집 잔여 {Remaining}페이지 (정렬 desc 라 가장 오래된 공시부터 유실). 범위=[{Begin}~{End}] 를 좁혀 재수집 필요", JobName, totalPage, MaxPages, totalPage - MaxPages, bgnDe, endDe); } return (true, inserted, skipped); } }