8.1 KiB
8.1 KiB
QuantEngine 수집 파이프라인 (KIS API)
1. 수집 실행 상태 전이도 (State Diagram)
KIS 데이터 수집 실행(kis_collection_runs)의 상태 흐름. 상태값은 KisDataCollectionOrchestrator 에서 정의:
RUNNING: 수집 진행 중COMPLETED: 모든 스냅샷 수집 완료 (에러 없음,total_errors == 0)COMPLETED_WITH_ERRORS: 부분 수집 완료 (에러 발생,total_errors > 0이지만 일부 성공)FAILED: 전체 실패 (예외 발생, 데이터 미적재)
stateDiagram-v2
[*] --> RUNNING: 수집 시작<br/>(RunCollectionAsync)
RUNNING --> COMPLETED: 완료 & error_count==0
RUNNING --> COMPLETED_WITH_ERRORS: 완료 & error_count>0
RUNNING --> FAILED: 예외 발생
COMPLETED --> [*]
COMPLETED_WITH_ERRORS --> [*]
FAILED --> [*]
상태 전이 조건 (KisDataCollectionOrchestrator.cs 라인 104-105):
error_count == 0→COMPLETEDerror_count > 0→COMPLETED_WITH_ERRORS- 예외(Exception) →
FAILED
성공 기준 (CLAUDE.md "Collection Run Success Criteria"):
- Success:
status == "COMPLETED"(NOT failed) - Partial Success:
status == "COMPLETED"+total_snapshots > 0+total_errors > 0 - Failure:
status == "FAILED"ORtotal_snapshots == 0
2. 수집 파이프라인 흐름도 (Flowchart)
KIS API 데이터 수집의 전체 흐름. 두 개의 트리거:
- Hangfire 정기 작업: 매일 09:00 에 자동 실행
- API 수동 트리거: POST /api/collection/run (쿠키 기반 인증)
flowchart TD
A["Hangfire daily-collection<br/>(09:00 KST)"]
B["POST /api/collection/run<br/>(Cookie Auth)"]
A --> C["IServiceScopeFactory.CreateScope<br/>(resolve ICollectionOrchestrator)"]
B --> C
C --> D["KisDataCollectionOrchestrator.RunCollectionAsync<br/>(tickers: [005930, 000660, ...])"]
D --> E["Per-ticker 루프"]
E --> F["KisApiPriceSource.GetPriceDataAsync<br/>(ticker, account)"]
F --> G["PriceDataNormalizer.NormalizeCollectionRow<br/>(seedRow, kisResult)"]
G --> H["CollectionRepository.SaveSnapshot<br/>(kis_collection_snapshots)"]
G --> I["CollectionRepository.SaveError<br/>(kis_collection_errors, on exception)"]
H --> J{루프 끝?}
I --> J
J -->|Yes| K["CollectionRepository.SaveRun<br/>(kis_collection_runs)"]
J -->|No| E
K --> L["파일 출력:<br/>Temp/kis_dotnet_collection_v1.json"]
L --> M["Serilog 로그:<br/>src/dotnet/.../logs/"]
M --> N["Admin UI: /Admin/Collection<br/>(CollectionRepository 읽기)"]
N --> O["대시보드 표시:<br/>상태, 스냅샷 수, 에러"]
데이터 흐름:
- 입력: Hangfire 스케줄 or API 수동 요청
- 오케스트레이션: ICollectionOrchestrator 스코프 생성
- 수집: KIS Open API 호출 → PriceDataNormalizer → DB 저장
- 출력:
- kis_collection_runs: 실행 메타데이터 (run_id, status, total_snapshots, total_errors)
- kis_collection_snapshots: 종목별 가격 데이터 (JSON payload)
- kis_collection_errors: 에러 기록
- Temp/kis_dotnet_collection_v1.json: 수집 결과 요약 (formula_id, gate, run_id, summary)
- Serilog 로그: 런타임 로그 (src/dotnet/QuantEngine.Web/logs/)
- 표시: Admin UI에서 CollectionRepository API 호출 → kis_collection_* 읽기 → Dashboard 렌더링
3. WBS 증거 검증 시퀀스도 (Sequence Diagram)
작업 완료 증거를 자동 검증하는 파이프라인. 도구: verify_wbs_task_v1.py (증거 수집) + validate_quant_engine_wbs_v1.py (CI에서 재검증).
sequenceDiagram
Developer->>verify_wbs_task_v1.py: python verify_wbs_task_v1.py --task QE-M1-01<br/>(또는 --run-commands)
verify_wbs_task_v1.py->>+spec/60_quant_engine_wbs.yaml: load spec
spec/60_quant_engine_wbs.yaml-->>-verify_wbs_task_v1.py: meta + tasks[QE-M1-01]
Note over verify_wbs_task_v1.py: evidence_checks 선언형 해석
alt pg_query 체크
verify_wbs_task_v1.py->>+PostgreSQL: SELECT ... (WHERE 절)
PostgreSQL-->>-verify_wbs_task_v1.py: 스칼라 결과 또는 행
verify_wbs_task_v1.py->>verify_wbs_task_v1.py: expect{min,max,equals} 비교
end
alt log_pattern 체크
verify_wbs_task_v1.py->>+src/dotnet/.../logs/: file_glob 매칭
src/dotnet/.../logs/-->>-verify_wbs_task_v1.py: 로그 라인
verify_wbs_task_v1.py->>verify_wbs_task_v1.py: 정규식 패턴 검사<br/>(min_matches, max_age_hours)
end
alt json_gate 체크
verify_wbs_task_v1.py->>+Temp/kis_dotnet_collection_v1.json: read JSON
Temp/kis_dotnet_collection_v1.json-->>-verify_wbs_task_v1.py: payload
verify_wbs_task_v1.py->>verify_wbs_task_v1.py: 점 표기 경로(dot notation)<br/>+ 값 비교 (>=N 지원)
end
alt file_exists 체크
verify_wbs_task_v1.py->>verify_wbs_task_v1.py: paths[] 존재 확인<br/>(min_bytes 검증)
end
alt playwright_report 체크
verify_wbs_task_v1.py->>+tests/e2e/playwright-report.json: read report
tests/e2e/playwright-report.json-->>-verify_wbs_task_v1.py: suites[].specs[]
verify_wbs_task_v1.py->>verify_wbs_task_v1.py: spec_file 매칭<br/>(passed_min, failed)
end
verify_wbs_task_v1.py->>verify_wbs_task_v1.py: 모든 체크 결과 종합<br/>(gate = ALL PASS? → PASS : FAIL)
verify_wbs_task_v1.py->>+Temp/evidence/QE-M1-01/: mkdir
verify_wbs_task_v1.py->>Temp/evidence/QE-M1-01/verdict.json: write verdict<br/>(task_id, gate, checks[])
verify_wbs_task_v1.py->>Temp/evidence/QE-M1-01/: save raw evidence<br/>(pg_query_n.json, log_excerpt.txt, ...)
verify_wbs_task_v1.py->>+runtime/lineage_events.jsonl: append event<br/>(node_id, gate, timestamp)
Developer<<--verify_wbs_task_v1.py: exit 0 (gate=PASS)<br/>or exit 1 (gate=FAIL)
Note over Developer: 선택: --run-commands 플래그<br/>verification_commands[] 실행
Developer->>+validate_quant_engine_wbs_v1.py: (CI) python validate_quant_engine_wbs_v1.py
validate_quant_engine_wbs_v1.py->>validate_quant_engine_wbs_v1.py: spec load
validate_quant_engine_wbs_v1.py->>validate_quant_engine_wbs_v1.py: tasks[status==DONE] 필터
validate_quant_engine_wbs_v1.py->>+Temp/evidence/*/verdict.json: load all verdicts
Temp/evidence/*/verdict.json-->>-validate_quant_engine_wbs_v1.py: gate 값
validate_quant_engine_wbs_v1.py->>validate_quant_engine_wbs_v1.py: gate=FAIL? → CI FAIL
validate_quant_engine_wbs_v1.py->>+Temp/quant_engine_wbs_v1.json: write summary
Developer<<--validate_quant_engine_wbs_v1.py: exit 0 (모두 PASS)<br/>or exit 1 (일부 FAIL)
검증 프로세스 상세:
| 단계 | 역할 | 산출물 |
|---|---|---|
| 1. 스펙 로드 | verify_wbs_task_v1.py | spec/60_quant_engine_wbs.yaml |
| 2. 증거 체크 실행 | 선언형 evidence_checks[] | pg_query / log_pattern / json_gate / file_exists / playwright_report |
| 3. 게이트 결정 | 모든 체크 PASS? | gate = PASS or FAIL |
| 4. 증거 저장 | Temp/evidence/<TASK_ID>/ | verdict.json + 원시 증거 |
| 5. 계보 로깅 | runtime/lineage_events.jsonl | node_id, gate, timestamp |
| 6. CI 재검증 | validate_quant_engine_wbs_v1.py | status=DONE 작업만 재검증 |
주요 특징:
- 선언형 검증: 체크 로직을 YAML에 기술 (하드코딩 최소화)
- 원시 증거 보존: 각 체크의 상세 결과를 JSON/텍스트로 저장
- 완료 주장 차단: "완료했다"는 수동 선언 불가 → verdict.json gate=PASS만 인정
- CI 편입: validate_quant_engine_wbs_v1.py가 release DAG의 노드로 동작
- 멀티 트리거: 단일 작업 검증 (--task) 또는 전체 검증 (CI)
검증 체크 타입 참고 (spec/60_quant_engine_wbs.yaml "evidence_check_types"):
- pg_query: PostgreSQL 스칼라 결과 비교 (min/max/equals)
- log_pattern: 로그 파일 정규식 매칭 (min_matches, max_age_hours)
- json_gate: JSON 아티팩트 키-값 검사 (점 표기 경로, >=N 비교)
- file_exists: 파일 존재 + 크기 검증 (min_bytes)
- playwright_report: Playwright 리포트 테스트 결과 (passed_min, failed)