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;
if (string.IsNullOrWhiteSpace(cfg.ApiKey))
{
Logger.LogWarning("[{Job}] OpenDart:ApiKey 미설정 — 수집 skip", JobName);
return true;
}
using var scope = scopeFactory.CreateScope();
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);
}
}