diff --git a/docs/design/Gimini-3-#92-worker-polling-indexing-orchestration.md b/docs/design/Gimini-3-#92-worker-polling-indexing-orchestration.md new file mode 100644 index 0000000..40b8dfa --- /dev/null +++ b/docs/design/Gimini-3-#92-worker-polling-indexing-orchestration.md @@ -0,0 +1,427 @@ +# Issue #92 Worker 자동 Polling 및 인덱싱 실행 오케스트레이션 상세 설계 + +closes #92 + +## 1. 문서 목적 + +이 문서는 이슈 [#92](https://github.com/DocGrid/backend/issues/92)의 구현 기준을 정의한다. + +현재 인덱싱 도메인은 Worker 등록·Heartbeat, Job Claim, Attempt 시작, 문서 Chunk 생성, Embedding 생성, +인덱싱 완료, 협력적 실패 보고, Lease 갱신과 만료 복구를 각각 제공한다. 그러나 각 기능은 관리자 API나 +Service 호출 단위로만 연결돼 있어 Worker를 활성화해도 PENDING Job을 자동으로 가져와 끝까지 실행하지 +않는다. + +이번 작업은 같은 Spring 애플리케이션 안에서 실행되는 Worker가 사용 가능한 로컬 실행 슬롯만큼 Job을 +Claim하고, 기존 Service를 직접 조합해 파이프라인을 수행하도록 한다. 외부 파일 저장소와 Embedding +Provider 호출 중에는 DB Transaction을 유지하지 않으며, 현재 Worker ID와 Claim Token은 각 단계의 +소유권 검증에만 사용한다. + +### 1.1 성공 기준 + +- Worker 기능이 비활성화된 API 전용 실행에서는 Poller와 실행 Executor가 만들어지지 않는다. +- Worker 등록이 완료된 뒤에만 Job을 Polling한다. +- `max-concurrency`를 초과해 Job을 Claim하거나 실행하지 않는다. +- 실행 슬롯을 먼저 확보한 뒤 한 Job만 Claim하고, Job이 없거나 제출에 실패하면 슬롯을 즉시 반환한다. +- Claim한 Job은 Attempt 시작, Chunk 생성, Embedding 생성, 완료 순서로 처리한다. +- 이미 진행된 문서 버전은 현재 상태에 맞는 단계부터 안전하게 재개한다. +- 외부 I/O 중 DB 행 잠금이나 장기 Transaction을 유지하지 않는다. +- 실행 중인 Job은 설정된 주기로 Lease를 갱신하고 종료 시 갱신 작업을 해제한다. +- Attempt 시작 이후의 실행 오류는 제한된 실패 유형과 안전한 메시지로 실패 Service에 보고한다. +- 소유권·Lease를 이미 잃은 실행은 과거 Claim으로 실패 상태를 덮어쓰지 않는다. +- 애플리케이션 종료 시 신규 Polling을 먼저 중단하고, 제한 시간 동안 실행 중인 작업을 기다린 뒤 Worker를 + STOPPED 처리한다. +- Claim Token, 원문 내용, Vector와 외부 인증 정보는 일반 로그에 남기지 않는다. +- 단위 테스트와 실제 PostgreSQL 동시성 테스트로 실행 상한, Lease, 종료와 단일 Claim을 검증한다. + +## 2. 범위 + +### 2.1 포함 범위 + +- Polling 주기, 최대 동시 실행 수, Lease 갱신 주기와 종료 유예 시간 설정 +- 고정 크기·무대기 Job Executor와 단일 Lease Scheduler +- 로컬 실행 슬롯 예약·반환 +- Worker 등록 완료 여부를 기준으로 한 Polling +- 기존 Claim Service의 반복 호출 +- 문서 버전 상태 조회와 단계 재개 결정 +- Attempt, Chunk, Embedding, 완료 Service 직접 오케스트레이션 +- 실행별 Lease 갱신 등록·해제와 소유권 상실 표시 +- 예외별 실패 유형 분류와 안전한 오류 메시지 생성 +- Attempt 시작 이후의 협력적 실패 보고 +- 신규 Claim 차단, 실행 대기, 강제 중단 순서의 Graceful Shutdown +- 단위, Spring Context, 실제 PostgreSQL과 동시성 검증 + +### 2.2 제외 범위 + +- Message Broker, 외부 작업 Queue와 분산 Lock +- Claim한 로컬 작업의 영속 Queue +- Worker 수에 따른 자동 확장 +- Chunk·Embedding 단위 Checkpoint와 부분 Batch 재개 +- Retry Jitter와 관리자 수동 재시도·취소 API +- Worker 전용 Machine Credential +- 실행 Dashboard, Metric, Alert와 분산 Trace +- 외부 Worker 프로세스를 위한 HTTP Client +- 파일 형식과 Embedding 모델 정책 변경 + +이번 Worker는 기존 Service들과 같은 Spring Context에서 동작한다. 자기 자신에게 관리자 HTTP 요청을 +보내면 인증·직렬화·네트워크 실패 지점만 추가되므로 내부 Service를 직접 호출한다. 외부 프로세스 Worker가 +필요해지는 시점에 별도 Client 계약을 설계한다. + +## 3. 현재 기준선 + +### 3.1 Worker 생명주기 + +`WorkerLifecycleManager`는 애플리케이션 준비 이벤트에서 Worker를 등록하고 현재 Worker ID를 메모리에 +보관한다. Heartbeat Scheduler는 Worker ID가 존재할 때만 갱신한다. Context 종료 이벤트에서는 Worker를 +STOPPED로 변경하고 메모리의 ID를 제거한다. + +Poller는 이 ID를 등록 완료 신호로 사용한다. 등록 전에는 Claim Service를 호출하지 않으며 등록이 실패해 +ID가 없으면 다음 주기에도 아무 작업을 하지 않는다. + +### 3.2 Claim과 실행 소유권 + +Claim Service는 PostgreSQL `FOR UPDATE SKIP LOCKED`로 다음 PENDING Job을 하나 선택하고 아래 값을 한 +Transaction에서 기록한다. + +~~~text +status = PROCESSING +locked_by_worker_id = 현재 Worker +claim_token = 새 UUID +locked_at = Claim 시각 +lock_expires_at = Claim 시각 + leaseDuration +~~~ + +각 후속 Service는 Job 행을 먼저 잠그고 Worker ID, Claim Token과 유효 Lease를 다시 검증한다. 로컬 실행 +상태는 편의를 위한 조정 정보일 뿐이며, 결과 저장의 최종 권한은 DB 소유권 검증이 결정한다. + +### 3.3 단계별 Transaction 경계 + +Chunk와 Embedding Service는 준비 Transaction과 완료 Transaction 사이에서 외부 I/O를 수행한다. + +~~~text +짧은 준비 Transaction +→ Transaction 밖 파일 읽기·파싱 또는 Embedding HTTP 호출 +→ 짧은 완료 Transaction +~~~ + +오케스트레이터에는 `@Transactional`을 적용하지 않는다. 기존 Service가 가진 짧은 Transaction 경계를 +그대로 사용해 Object Storage와 Embedding Provider 응답을 기다리는 동안 DB 잠금을 유지하지 않는다. + +### 3.4 단계 재생 범위 + +- Chunk Service는 `UPLOADED`, `PARSING`에서 작업하고 `CHUNKED` 결과를 재생한다. +- Embedding Service는 `CHUNKED`, `EMBEDDING`에서 작업 또는 결과를 재생한다. +- 완료 Service는 `EMBEDDING`의 완전한 결과를 `INDEXED`로 확정하고 같은 실행의 완료를 재생한다. + +Chunk Service는 `EMBEDDING` 상태를 재생하지 않는다. 따라서 모든 실행에서 Chunk를 무조건 다시 호출하지 +않고 현재 문서 버전 상태를 조회해 시작 단계를 선택해야 한다. + +## 4. 핵심 결정 + +### 4.1 실행 슬롯을 Claim보다 먼저 확보한다 + +동시 실행 상한은 Executor Queue 크기가 아니라 명시적인 슬롯으로 관리한다. + +~~~text +1. 로컬 슬롯 확보 +2. 등록된 Worker ID 확인 +3. Job 하나 Claim +4. Executor에 즉시 제출 +5. 작업 종료 또는 중간 실패 시 슬롯 반환 +~~~ + +슬롯이 없으면 Claim Service를 호출하지 않는다. Claim 후 Queue에서 오래 기다리며 Lease를 소모하는 상황을 +막기 위해 Job Executor는 고정 Thread 수와 `SynchronousQueue`를 사용한다. 즉시 실행할 Thread가 없으면 +제출을 거부하고 슬롯을 반환한다. 정상 설계에서는 슬롯 수와 Thread 수가 같으므로 제출 거부는 종료 경쟁이나 +내부 불변식 오류일 때만 발생한다. + +한 Polling 주기에는 실행 가능한 슬롯 수만큼 위 절차를 반복한다. 첫 빈 Claim을 만나면 현재 PENDING +후보가 없다고 판단해 그 주기의 반복을 끝낸다. + +### 4.2 단계는 문서 버전 상태로 재개한다 + +Attempt를 시작한 뒤 Claim 응답의 `documentVersionId`로 현재 상태를 조회한다. + +| 현재 상태 | 실행 | +| --- | --- | +| UPLOADED | Chunk 생성 → Embedding 생성 → 완료 | +| PARSING | Chunk 재개 → Embedding 생성 → 완료 | +| CHUNKED | Embedding 생성 → 완료 | +| EMBEDDING | Embedding 재개 → 완료 | +| INDEXED | PROCESSING Job과 모순이므로 실패 보고 | +| FAILED | 실행 대상이 아니므로 실패 보고 | + +상태 조회는 경로 선택용 Snapshot일 뿐이다. 조회 직후 상태가 바뀔 수 있으므로 각 단계 Service의 Job 잠금, +Attempt와 소유권 검증이 최종 정확성을 보장한다. 오케스트레이터가 JPA Entity를 단계 사이에 보관하지 +않는다. + +### 4.3 Claim과 Attempt 사이 오류는 합성 실패를 만들지 않는다 + +Claim 성공 후 Attempt 시작 전에 프로세스가 종료되거나 Attempt Service가 실패할 수 있다. 이때 실제 +Attempt ID가 없으므로 실패 Service를 호출하거나 가짜 Attempt를 만들지 않는다. Lease 갱신도 시작하지 +않고 소유권 만료 복구가 Job을 회수하게 둔다. + +Attempt가 시작된 뒤 발생한 오류만 현재 Attempt ID로 협력적 실패를 보고한다. + +### 4.4 실패 분류는 제한된 계약만 저장한다 + +예외 메시지와 Stack Trace를 그대로 DB에 저장하면 원문, Object Key, 외부 Endpoint나 인증 정보가 섞일 수 +있다. 분류기는 `DocGridException.errorCode`를 허용 목록으로 매핑하고 고정된 안전 메시지를 만든다. + +| 원인 | 실패 유형 | Retry | +| --- | --- | --- | +| FILE_STORAGE_FAILED | STORAGE_UNAVAILABLE | 가능 | +| 지원하지 않는 형식, 빈 내용, UTF-8 해석 실패 | DOCUMENT_CONTENT_INVALID | 불가 | +| EMBEDDING_SERVER_UNAVAILABLE | EMBEDDING_PROVIDER_UNAVAILABLE | 가능 | +| 차원 불일치, 잘못된 Vector | EMBEDDING_RESULT_INVALID | 불가 | +| 상태·연관·Chunk·Embedding 불변식 오류 | INDEXING_STATE_INCONSISTENT | 불가 | +| 그 밖의 실행 오류 | WORKER_INTERNAL_ERROR | 가능 | + +다음 오류는 이미 현재 실행의 권한을 잃었음을 뜻하므로 실패 보고를 시도하지 않는다. + +- EMBEDDING_JOB_NOT_FOUND +- EMBEDDING_JOB_NOT_PROCESSING +- EMBEDDING_JOB_OWNERSHIP_INVALID +- EMBEDDING_JOB_LEASE_EXPIRED +- EMBEDDING_JOB_ATTEMPT_INVALID +- EMBEDDING_JOB_FAILURE_CONFLICT + +이 경우 현재 실행은 중단하고 Lease 복구 또는 이미 완료된 경쟁 실행의 결과를 따른다. 실패 보고 자체가 +실패해도 원래 예외를 숨기지 않으며 Claim Token 없이 Job ID, Attempt ID와 오류 코드만 로그에 남긴다. + +### 4.5 Lease 갱신은 실행별 Handle로 관리한다 + +Attempt 시작 직후 실행별 Lease 갱신 Handle을 만든다. 하나의 Scheduled Executor가 각 활성 실행의 갱신을 +예약하며 갱신 요청은 기존 `EmbeddingJobLeaseService`를 직접 호출한다. + +~~~text +Attempt 시작 +→ 즉시 Lease Handle 등록 +→ lease-renewal-interval마다 갱신 +→ 각 단계 전후 현재 소유권 확인 +→ 완료·실패·중단의 finally에서 Handle 해제 +~~~ + +갱신이 소유권·상태 충돌로 거부되면 Handle을 `lost` 상태로 바꾸고 이후 갱신을 취소한다. 파이프라인은 단계 +사이에서 이를 확인해 다음 외부 작업을 시작하지 않는다. 일반 인프라 예외는 로그에 기록하고 다음 예약을 +유지한다. 실제 결과 저장 전에는 기존 Service가 DB Lease를 다시 검증하므로 갱신 실패가 소유권을 연장한 +것처럼 취급되지 않는다. + +갱신 주기는 0보다 크고 Lease 기간보다 짧아야 한다. 기본값은 Lease 5분, 갱신 1분이다. + +### 4.6 종료는 Polling 중단 후 Worker 정지 순서다 + +Context 종료 이벤트 Listener 순서를 명시한다. + +~~~text +1. Poller가 신규 슬롯 확보와 Claim 중단 +2. Job Executor shutdown +3. shutdown-grace-period 동안 활성 실행 완료 대기 +4. 시간 초과 시 실행 Thread interrupt와 남은 Lease Handle 취소 +5. WorkerLifecycleManager가 Worker를 STOPPED 처리 +6. 완료되지 않은 PROCESSING Job은 Lease 만료 복구가 회수 +~~~ + +실행 대기 중 Worker ID와 Heartbeat를 유지해야 진행 중인 작업이 완료·실패와 Lease 갱신을 수행할 수 있다. +따라서 실행 종료 Listener를 높은 우선순위로, Worker STOPPED Listener를 낮은 우선순위로 둔다. 유예 시간이 +끝난 뒤 DB 상태를 임의로 실패 처리하지 않는다. 외부 I/O가 interrupt에 반응하지 않더라도 만료된 Lease가 +과거 결과 저장을 차단한다. + +## 5. 구성 요소 설계 + +### 5.1 설정과 실행 기반 + +`IndexingWorkerProperties`에 다음 값을 추가한다. + +| 설정 | 기본값 | 검증 | +| --- | --- | --- | +| polling-interval | 1초 | 양수 | +| max-concurrency | 2 | 1 이상 | +| lease-renewal-interval | 1분 | 양수, lease-duration보다 짧음 | +| shutdown-grace-period | 30초 | 0 이상 | + +종료 유예 시간은 0을 허용해 즉시 중단 정책을 표현한다. 나머지 주기는 0 또는 음수를 허용하지 않는다. + +`WorkerExecutionConfig`는 Worker 활성화 조건에서만 다음 Bean을 제공한다. + +- `workerJobExecutor`: `max-concurrency` 고정 Thread, `SynchronousQueue`, AbortPolicy +- `workerLeaseScheduler`: 단일 Scheduled Thread, 취소 작업 즉시 제거 + +Thread 이름에는 역할과 번호만 포함하고 Worker 이름, Job ID와 Claim Token을 포함하지 않는다. + +### 5.2 실행 슬롯 + +`WorkerExecutionSlotPool`은 `Semaphore(max-concurrency)`와 신규 작업 허용 상태를 소유한다. 획득 성공 시 +한 번만 닫을 수 있는 `WorkerExecutionSlot`을 반환한다. Claim 없음, 제출 거부, 파이프라인 종료가 같은 +슬롯을 중복 반환하지 않도록 Slot 자체가 원자적인 closed 상태를 가진다. + +### 5.3 상태 조회 + +`DocumentIndexingStageQueryService`는 `documentVersionId`로 `DocumentVersionStatus`만 반환한다. 조회 결과는 +분기용 Snapshot이며 Entity를 Worker 계층에 노출하지 않는다. Version이 없으면 기존 문서 상태 불변식 오류로 +처리한다. + +### 5.4 파이프라인 + +`WorkerIndexingPipeline`은 다음 의존성을 조합한다. + +- EmbeddingJobAttemptService +- DocumentIndexingStageQueryService +- DocumentParsingService +- DocumentEmbeddingService +- DocumentIndexingCompletionService +- WorkerIndexingFailureReporter +- WorkerLeaseRenewalManager + +실행 순서는 다음과 같다. + +~~~text +1. Attempt 시작 +2. Lease 갱신 Handle 등록 +3. 현재 Version 상태 조회 +4. 필요한 경우 Chunk 생성 또는 재개 +5. 필요한 경우 Embedding 생성 또는 재개 +6. 인덱싱 완료 +7. 오류 시 실패 분류·보고 +8. finally에서 Lease Handle과 실행 슬롯 해제 +~~~ + +Pipeline 메서드에는 Transaction을 적용하지 않는다. Claim Token은 요청 DTO 생성에만 전달하고 로그나 결과 +객체의 `toString()` 출력에 포함하지 않는다. + +### 5.5 Poller와 실행 관리자 + +`WorkerJobPollingScheduler`는 고정 지연으로 `poll()`을 호출한다. 동시 Scheduler 실행을 막기 위해 현재 +Polling 여부를 원자적으로 보호한다. `poll()`은 Worker ID와 남은 슬롯을 확인하고, 확보한 슬롯마다 Claim과 +제출을 한 번 수행한다. + +`WorkerExecutionLifecycleManager`는 Poller 중단과 Executor 종료를 조정한다. 종료 Listener는 여러 번 +호출돼도 같은 종료 절차를 반복하지 않는다. InterruptedException을 받으면 현재 Thread의 interrupt 상태를 +복구하고 즉시 강제 중단 단계로 이동한다. + +## 6. 동시성·오류 계약 + +### 6.1 여러 Worker의 Polling + +각 Worker가 동시에 Polling해도 Claim Repository의 `FOR UPDATE SKIP LOCKED`가 같은 PENDING Job의 중복 +선택을 막는다. 로컬 슬롯은 한 프로세스의 실행 상한만 담당하고 전역 동시성 제어로 사용하지 않는다. + +### 6.2 종료와 Polling 경쟁 + +종료 플래그를 변경한 뒤 이미 슬롯을 확보한 Polling Thread가 있을 수 있다. Claim 직전에 허용 상태를 다시 +검사하고, Claim 이후 제출이 거부되면 슬롯을 반환한다. 이미 Claim된 Job은 즉시 실행할 수 없으면 Lease +만료 복구 대상으로 남긴다. 소유권을 임의로 반납하는 새 DB 전이는 만들지 않는다. + +### 6.3 Lease 갱신과 완료·실패 경쟁 + +갱신, 완료와 실패는 모두 Job 행을 먼저 잠근다. 먼저 Commit한 상태가 후속 호출의 소유권 검증 결과를 +결정한다. + +~~~text +갱신 선행 → 새 Lease 안에서 완료·실패 가능 +완료 선행 → 갱신은 PROCESSING 아님으로 거부 +실패 선행 → 갱신은 PROCESSING 아님 또는 소유권 제거로 거부 +복구 선행 → 과거 실행의 후속 결과 저장 거부 +~~~ + +### 6.4 Polling 실패 + +- Worker 미등록: 조용히 Skip +- 슬롯 없음: 조용히 Skip +- PENDING Job 없음: 정상 종료 +- Claim DB 오류: 해당 주기 중단, 다음 주기 재시도 +- Executor 제출 거부: 슬롯 반환, Job은 Lease 복구 대기 +- 파이프라인 RuntimeException: Attempt 존재 시 제한된 실패 보고 + +반복 Scheduler 자체가 예외로 중단되지 않도록 Polling 경계에서 RuntimeException을 기록하고 삼킨다. + +## 7. 보안·관측성 + +### 7.1 로그 허용 필드 + +- workerId +- jobId +- attemptId +- 문서 버전 상태 +- 실패 유형 또는 ErrorCode +- 실행 결과와 소요 시간 +- 활성 실행 수와 종료 대기 결과 + +### 7.2 로그 금지 필드 + +- Claim Token +- 문서 원문과 Chunk Text +- Embedding Vector +- Object Storage Bucket/Object Key +- 외부 Provider 요청·응답 본문 +- Authorization Header와 환경 변수 값 + +예상된 빈 Polling은 로그를 남기지 않는다. Job 시작·완료는 INFO, 소유권 상실과 제출 거부는 WARN, +불변식 오류와 실패 보고 실패는 ERROR를 사용한다. + +## 8. 테스트 전략 + +### 8.1 설정·슬롯 단위 테스트 + +- 새 설정의 기본값 +- 0·음수 Polling/갱신 주기 거부 +- Lease 이상 갱신 주기 거부 +- 0 미만 종료 유예 시간 거부 +- 최대 슬롯까지만 획득 +- Slot 중복 close가 Permit을 중복 반환하지 않음 +- Poller 중단 뒤 신규 슬롯 획득 거부 + +### 8.2 파이프라인 단위 테스트 + +- UPLOADED/PARSING은 Chunk부터 실행 +- CHUNKED/EMBEDDING은 Embedding부터 실행 +- 완료까지 Service 호출 순서 보장 +- Attempt 시작 전 오류는 실패 Service 미호출 +- Attempt 시작 후 오류는 분류된 실패 요청으로 보고 +- Claim Token이 로그·응답용 객체에 노출되지 않음 +- 소유권 상실 오류는 실패 보고 미호출 +- 모든 종료 경로에서 Lease Handle과 슬롯 반환 + +### 8.3 Lease 단위 테스트 + +- 설정 주기로 갱신 Service 호출 +- Handle 종료 시 예약 작업 취소 +- 소유권 오류 시 lost 표시와 후속 예약 중단 +- 일시적 인프라 오류 뒤 예약 유지 +- 완료와 실패 경로에서 활성 Handle 제거 + +### 8.4 Polling·종료 단위 테스트 + +- Worker 등록 전 Claim 미호출 +- 가용 슬롯 수만큼만 Claim·제출 +- 빈 Claim에서 추가 조회 중단 +- Claim 예외와 제출 거부 시 슬롯 반환 +- 종료 후 신규 Claim 없음 +- 유예 시간 안 완료 시 강제 중단 없음 +- 유예 시간 초과 시 `shutdownNow`와 Lease Handle 취소 +- Worker STOPPED 기록이 실행 종료 대기 뒤 수행됨 + +### 8.5 실제 PostgreSQL 통합·동시성 테스트 + +- 여러 Poller가 경쟁해도 한 Job은 한 Worker만 Claim +- 한 Worker의 활성 PROCESSING Job 수가 설정된 최대 동시성 이하 +- 장기 실행 중 Lease 갱신으로 만료 복구가 Job을 회수하지 않음 +- 갱신 중단 후 만료되면 Recovery가 Job을 한 번만 재예약 +- 종료 경쟁에서 완료된 Job은 INDEXED, 미완료 Job은 Lease 만료 뒤 재시도 가능 +- 기존 Claim, Attempt, Chunk, Embedding, 완료와 Lease 복구 회귀 테스트 통과 + +실행된 명령, 환경, 결과와 측정값은 구현 완료 후 `docs/test-results/` 문서에 기록한다. + +## 9. 구현 순서 + +1. 이 상세 설계 문서 확정 +2. Worker 실행 설정, Executor와 슬롯 기반 추가 +3. 상태 조회와 인덱싱 Pipeline 구현 +4. 실패 분류와 보고 연결 +5. 실행별 Lease 갱신 관리 +6. 슬롯 기반 Polling과 Graceful Shutdown 연결 +7. 설정·슬롯·Pipeline·Lease·Polling 단위 테스트 +8. 실제 PostgreSQL 통합·동시성 테스트 +9. 전체 검증 결과 문서화 + +각 단계는 독립적으로 빌드 가능한 커밋으로 유지한다. 구현 중 설계 변경이 필요하면 먼저 이 문서의 관련 +결정과 테스트 기준을 갱신한다. diff --git a/docs/test-results/Gimini-3-#92-worker-polling-indexing-orchestration.md b/docs/test-results/Gimini-3-#92-worker-polling-indexing-orchestration.md new file mode 100644 index 0000000..333d172 --- /dev/null +++ b/docs/test-results/Gimini-3-#92-worker-polling-indexing-orchestration.md @@ -0,0 +1,116 @@ +# #92 Worker Polling 및 인덱싱 실행 오케스트레이션 검증 결과 + +## 1. 검증 정보 + +- 실행일: 2026-08-03 (Asia/Seoul) +- 대상 브랜치: `feature/92` +- 애플리케이션: Spring Boot 3.5.16, Java 17 +- 데이터베이스: Docker Desktop의 격리된 PostgreSQL 14.6(OpenSQL 호환) + pgvector +- 스키마: Flyway V1~V35 적용, 통합 테스트별 격리 스키마 사용 +- 최종 결과: 정규 테스트 576개와 동시성 테스트 10개 통과, 실패·오류·Skip 0개 + +검증에는 localhost에만 노출한 작업 전용 컨테이너와 데이터 볼륨을 사용했다. 기존 +`local-opensql` 컨테이너와 `opensql_data` 공유 볼륨은 변경하지 않았다. 검증 종료 후 작업 전용 +컨테이너와 볼륨은 삭제했으며, 운영 Secret은 사용하거나 기록하지 않았다. + +## 2. Swagger/OpenAPI 수동 검증 + +- 결과: 해당 없음 +- 근거: 이번 변경은 Worker 내부 Scheduler, 실행 Service와 설정만 추가하며 Controller, 요청·응답 DTO, + Endpoint 및 OpenAPI 계약을 변경하지 않는다. + +따라서 Swagger에서 호출할 신규·변경 API가 없으며, Worker 내부 실행 계약은 아래 자동 테스트와 실제 +PostgreSQL 동시성 테스트로 검증했다. + +## 3. 전체 정규 테스트 + +실행 명령의 환경 값은 Placeholder로 대체한다. + +```bash +DB_HOST=localhost \ +DB_PORT='' \ +DB_NAME=docgrid \ +DB_USER='' \ +DB_PASSWORD='' \ +DB_SSLMODE=disable \ +JWT_SECRET='' \ +./gradlew test +``` + +결과: + +```text +BUILD SUCCESSFUL +test suites=87 tests=576 failures=0 errors=0 skipped=0 +``` + +기본 `test` Task에서 Worker 등록·Heartbeat·상태 관리, Polling Scheduler, 실행 슬롯, 파이프라인, +Lease 갱신, 실패 보고, 종료 절차와 기존 도메인 회귀 테스트를 함께 검증했다. + +## 4. PostgreSQL 동시성 검증 + +실행: + +```bash +DB_HOST=localhost \ +DB_PORT='' \ +DB_NAME=docgrid \ +DB_USER='' \ +DB_PASSWORD='' \ +DB_SSLMODE=disable \ +JWT_SECRET='' \ +./gradlew claimConcurrencyTest +``` + +결과: + +```text +BUILD SUCCESSFUL +test suites=3 tests=10 failures=0 errors=0 skipped=0 +``` + +| 테스트 클래스 | 테스트 수 | 결과 | +| --- | ---: | --- | +| `EmbeddingJobClaimConcurrencyIntegrationTest` | 2 | 통과 | +| `EmbeddingJobLeaseRecoveryIntegrationTest` | 4 | 통과 | +| `WorkerOrchestrationIntegrationTest` | 4 | 통과 | + +새 Worker 오케스트레이션 통합 테스트는 다음 경쟁 조건을 실제 PostgreSQL 잠금과 트랜잭션으로 +검증했다. + +| 검증 항목 | 결과 | +| --- | --- | +| 두 Poller가 한 Job에 경쟁할 때 Claim과 파이프라인 제출이 한 번만 발생 | 통과 | +| 실행 슬롯 2개인 Worker가 5개 Job 중 2개만 PROCESSING으로 Claim | 통과 | +| 활성 실행의 Lease 갱신이 원래 만료 시각의 복구를 차단 | 통과 | +| Lease 갱신 중단 후 동시 복구가 Job을 정확히 한 번만 재예약 | 통과 | + +## 5. 표준 빌드 + +동일한 격리 DB와 테스트 전용 인증 설정에서 다음 명령을 실행했다. + +```bash +./gradlew build +``` + +결과: `BUILD SUCCESSFUL`. Compile, Test, Check, Boot JAR 및 JAR 생성 단계가 모두 성공했다. + +## 6. 확인된 불변식 + +- 여러 Worker가 같은 대기 Job을 조회해도 PostgreSQL Claim은 한 Worker에만 귀속된다. +- 한 Worker가 소유하는 PROCESSING Job 수는 로컬 실행 슬롯 수를 넘지 않는다. +- 슬롯이 없을 때 Poller는 추가 Job을 Claim하지 않고 다음 주기를 기다린다. +- 활성 파이프라인은 완료 전까지 Lease를 갱신하며, 갱신된 Job은 이전 만료 시각에 복구되지 않는다. +- Lease 갱신이 멈춘 만료 Job에 여러 복구 실행이 경쟁해도 Retry 전이는 한 번만 적용된다. +- 성공·재시도·최종 실패 경로에서 Lease 갱신과 실행 슬롯은 정리된다. +- 종료 요청 후 새 Polling은 시작되지 않고, 대기 중인 실행은 제한 시간 정책에 따라 정리된다. + +## 7. 환경 진단 기록과 제한 사항 + +- 기존 공유 OpenSQL 볼륨은 이미지가 기대하는 내부 Role과 초기화 상태가 달라 사용할 수 없었다. + 공유 데이터를 수정하지 않고 별도 작업 전용 컨테이너와 볼륨으로 전환했다. +- 최초 일부 기존 통합 테스트 실행은 테스트용 JWT 환경 값 누락으로 Spring Context 구성 단계에서 + 실패했다. 테스트 전용 값을 제공한 재실행과 최종 전체 실행은 모두 통과했다. +- 프로젝트 기본 `test` Task가 제외하는 `claim-concurrency`는 전용 Task로 별도 실행했다. +- 외부 MinIO 의존 테스트, 실제 임베딩 서버 네트워크 동작과 Benchmark는 이번 검증 범위가 아니다. +- 테스트용 컨테이너와 볼륨은 종료 시 삭제했으므로 그 안의 데이터는 복구하지 않는다. diff --git a/src/main/java/com/opensource/docgrid/domain/document/service/query/DocumentIndexingStageQueryService.java b/src/main/java/com/opensource/docgrid/domain/document/service/query/DocumentIndexingStageQueryService.java new file mode 100644 index 0000000..ffbe3d6 --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/document/service/query/DocumentIndexingStageQueryService.java @@ -0,0 +1,34 @@ +package com.opensource.docgrid.domain.document.service.query; + +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +import com.opensource.docgrid.domain.document.enums.DocumentVersionStatus; +import com.opensource.docgrid.domain.document.repository.DocumentVersionRepository; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +import lombok.RequiredArgsConstructor; + +/** + * Worker 파이프라인의 시작 단계를 결정할 문서 버전 상태 Snapshot을 조회한다. + * + *

조회 결과는 경로 선택에만 사용하며 Entity를 Worker 계층에 전달하지 않는다. 실제 상태 변경 가능 + * 여부와 소유권은 각 Command Service가 Job과 Version을 잠근 뒤 다시 검증한다. + */ +@Service +@RequiredArgsConstructor +@Transactional(readOnly = true) +public class DocumentIndexingStageQueryService { + + private final DocumentVersionRepository documentVersionRepository; + + /** + * Claim 응답이 가리키는 문서 버전의 현재 파이프라인 상태를 반환한다. + */ + public DocumentVersionStatus getStatus(Long documentVersionId) { + return documentVersionRepository.findById(documentVersionId) + .map(documentVersion -> documentVersion.getStatus()) + .orElseThrow(() -> new DocGridException(ErrorCode.INDEXING_STATUS_INCONSISTENT)); + } +} diff --git a/src/main/java/com/opensource/docgrid/domain/worker/config/IndexingWorkerProperties.java b/src/main/java/com/opensource/docgrid/domain/worker/config/IndexingWorkerProperties.java index 459ecad..accf8ad 100644 --- a/src/main/java/com/opensource/docgrid/domain/worker/config/IndexingWorkerProperties.java +++ b/src/main/java/com/opensource/docgrid/domain/worker/config/IndexingWorkerProperties.java @@ -14,7 +14,7 @@ import lombok.Setter; /** - * 인덱싱 Worker의 실행 여부, Heartbeat, DEAD 판정, Job Lease와 만료 복구를 바인딩하는 설정 클래스. + * 인덱싱 Worker의 실행 여부, Polling, 동시 실행, Heartbeat와 Job Lease 생명주기를 바인딩하는 설정 클래스. * *

{@code indexing.worker} 환경 설정을 타입 안전한 {@link Duration}으로 제공하고, 애플리케이션 시작 * 단계에서 서로 모순되거나 0 이하인 시간 설정을 차단한다. @@ -38,10 +38,21 @@ public class IndexingWorkerProperties { @NotNull private Duration deadThreshold = Duration.ofSeconds(30); + // 등록된 Worker가 실행 슬롯을 확인하고 새 Job을 찾는 주기다. + @NotNull + private Duration pollingInterval = Duration.ofSeconds(1); + + @Min(1) + private int maxConcurrency = 2; + // Claim 후 Worker가 소유권을 유지하는 기본 시간이다. 만료 복구는 후속 처리에서 사용한다. @NotNull private Duration leaseDuration = Duration.ofMinutes(5); + // 활성 실행은 Lease 만료 전에 이 주기로 소유권을 갱신한다. + @NotNull + private Duration leaseRenewalInterval = Duration.ofMinutes(1); + // 만료 Lease 복구 작업의 실행 주기와 한 번에 조회할 최대 Job 수다. @NotNull private Duration leaseRecoveryInterval = Duration.ofSeconds(30); @@ -56,6 +67,10 @@ public class IndexingWorkerProperties { @NotNull private Duration retryMaxDelay = Duration.ofMinutes(5); + // 종료 시 신규 Claim을 막은 뒤 활성 실행이 스스로 끝나기를 기다리는 최대 시간이다. + @NotNull + private Duration shutdownGracePeriod = Duration.ofSeconds(30); + /** * Heartbeat가 양수이고 DEAD 기준보다 짧은지 검증한다. */ @@ -68,6 +83,16 @@ public boolean isTimingValid() { && deadThreshold.compareTo(heartbeatInterval) > 0; } + /** + * 빈 작업 조회가 Busy Loop가 되지 않도록 Polling 주기가 양수인지 검증한다. + */ + @AssertTrue(message = "Job Polling 주기는 0보다 커야 합니다.") + public boolean isPollingIntervalValid() { + return pollingInterval != null + && !pollingInterval.isZero() + && !pollingInterval.isNegative(); + } + /** * 발급 즉시 만료되는 Lease가 만들어지지 않도록 Lease 기간이 양수인지 검증한다. */ @@ -78,6 +103,18 @@ public boolean isLeaseDurationValid() { && !leaseDuration.isNegative(); } + /** + * 활성 Job이 만료 전에 갱신될 수 있도록 갱신 주기가 양수이고 Lease 기간보다 짧은지 검증한다. + */ + @AssertTrue(message = "Lease 갱신 주기는 0보다 크고 Lease 기간보다 짧아야 합니다.") + public boolean isLeaseRenewalIntervalValid() { + return leaseRenewalInterval != null + && leaseDuration != null + && !leaseRenewalInterval.isZero() + && !leaseRenewalInterval.isNegative() + && leaseRenewalInterval.compareTo(leaseDuration) < 0; + } + /** * 만료 Lease 복구 Scheduler가 과도하게 반복되지 않도록 실행 주기가 양수인지 검증한다. */ @@ -99,4 +136,13 @@ public boolean isRetryDelayValid() { && !retryInitialDelay.isNegative() && retryMaxDelay.compareTo(retryInitialDelay) >= 0; } + + /** + * 즉시 종료는 허용하되 음수 대기 시간은 Executor 종료 계약으로 사용할 수 없으므로 차단한다. + */ + @AssertTrue(message = "Worker 종료 유예 시간은 0보다 작을 수 없습니다.") + public boolean isShutdownGracePeriodValid() { + return shutdownGracePeriod != null + && !shutdownGracePeriod.isNegative(); + } } diff --git a/src/main/java/com/opensource/docgrid/domain/worker/config/WorkerExecutionConfig.java b/src/main/java/com/opensource/docgrid/domain/worker/config/WorkerExecutionConfig.java new file mode 100644 index 0000000..90814d8 --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/worker/config/WorkerExecutionConfig.java @@ -0,0 +1,68 @@ +package com.opensource.docgrid.domain.worker.config; + +import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.SynchronousQueue; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.scheduling.concurrent.CustomizableThreadFactory; + +import com.opensource.docgrid.domain.worker.execution.WorkerExecutionSlotPool; + +/** + * Worker Job 실행과 Lease 갱신에 사용하는 제한된 Thread 자원을 구성한다. + * + *

Job Executor는 Queue에 Claim을 쌓지 않고 설정된 동시성만 즉시 실행한다. Lease Scheduler는 활성 + * 실행의 짧은 갱신 호출만 담당하며, Worker가 비활성화된 API 전용 실행에는 어떤 Thread도 만들지 않는다. + */ +@Configuration +@ConditionalOnProperty(prefix = "indexing.worker", name = "enabled", havingValue = "true") +public class WorkerExecutionConfig { + + public static final String WORKER_JOB_EXECUTOR = "workerJobExecutor"; + public static final String WORKER_LEASE_SCHEDULER = "workerLeaseScheduler"; + + /** + * 최대 동시 실행 수와 같은 크기의 무대기 Job Executor를 만든다. + */ + @Bean(name = WORKER_JOB_EXECUTOR, destroyMethod = "shutdownNow") + public ThreadPoolExecutor workerJobExecutor(IndexingWorkerProperties properties) { + int maxConcurrency = properties.getMaxConcurrency(); + return new ThreadPoolExecutor( + maxConcurrency, + maxConcurrency, + 0L, + TimeUnit.MILLISECONDS, + new SynchronousQueue<>(), + new CustomizableThreadFactory("indexing-worker-job-"), + new ThreadPoolExecutor.AbortPolicy() + ); + } + + /** + * 모든 활성 실행의 Lease 갱신을 직렬로 예약하는 단일 Thread Scheduler를 만든다. + */ + @Bean(name = WORKER_LEASE_SCHEDULER, destroyMethod = "shutdownNow") + public ScheduledThreadPoolExecutor workerLeaseScheduler() { + ScheduledThreadPoolExecutor scheduler = new ScheduledThreadPoolExecutor( + 1, + new CustomizableThreadFactory("indexing-worker-lease-") + ); + // 취소된 실행별 갱신 작업이 Scheduler Queue에 남아 종료와 메모리 회수를 늦추지 않게 한다. + scheduler.setRemoveOnCancelPolicy(true); + scheduler.setExecuteExistingDelayedTasksAfterShutdownPolicy(false); + scheduler.setContinueExistingPeriodicTasksAfterShutdownPolicy(false); + return scheduler; + } + + /** + * Claim 전에 실행 가능 여부를 예약하는 프로세스 로컬 슬롯 풀을 만든다. + */ + @Bean + public WorkerExecutionSlotPool workerExecutionSlotPool(IndexingWorkerProperties properties) { + return new WorkerExecutionSlotPool(properties.getMaxConcurrency()); + } +} diff --git a/src/main/java/com/opensource/docgrid/domain/worker/execution/WorkerExecutionSlotPool.java b/src/main/java/com/opensource/docgrid/domain/worker/execution/WorkerExecutionSlotPool.java new file mode 100644 index 0000000..1e05957 --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/worker/execution/WorkerExecutionSlotPool.java @@ -0,0 +1,89 @@ +package com.opensource.docgrid.domain.worker.execution; + +import java.util.Optional; +import java.util.concurrent.Semaphore; +import java.util.concurrent.atomic.AtomicBoolean; + +/** + * 한 Worker 프로세스가 Claim할 수 있는 Job 수를 실제 실행 가능 슬롯으로 제한한다. + * + *

분산 소유권은 DB Claim이 담당하며 이 클래스는 로컬 최대 동시성만 관리한다. Poller는 슬롯을 먼저 + * 확보해야 Job을 Claim할 수 있고, 종료가 시작되면 남은 Permit과 관계없이 신규 획득을 거부한다. + */ +public class WorkerExecutionSlotPool { + + private final Semaphore permits; + private final int capacity; + private final AtomicBoolean accepting = new AtomicBoolean(true); + + public WorkerExecutionSlotPool(int capacity) { + if (capacity < 1) { + throw new IllegalArgumentException("Worker 실행 슬롯 수는 1 이상이어야 합니다."); + } + this.capacity = capacity; + this.permits = new Semaphore(capacity, true); + } + + /** + * 현재 신규 실행을 허용하고 Permit이 남아 있을 때만 한 슬롯을 예약한다. + */ + public Optional tryAcquire() { + if (!accepting.get() || !permits.tryAcquire()) { + return Optional.empty(); + } + + // Permit 획득과 종료 플래그 변경이 경쟁하면 Claim 전에 즉시 반환해 종료 이후 신규 작업을 막는다. + if (!accepting.get()) { + permits.release(); + return Optional.empty(); + } + return Optional.of(new WorkerExecutionSlot(this)); + } + + /** + * 종료 절차가 시작된 뒤 신규 슬롯 획득을 영구적으로 중단한다. + */ + public void stopAccepting() { + accepting.set(false); + } + + public boolean isAccepting() { + return accepting.get(); + } + + public int getCapacity() { + return capacity; + } + + public int getAvailableSlots() { + return permits.availablePermits(); + } + + public int getActiveSlots() { + return capacity - permits.availablePermits(); + } + + private void release() { + permits.release(); + } + + /** + * 한 번 획득한 Worker 실행 Permit을 정확히 한 번 반환하는 수명 Handle이다. + */ + public static final class WorkerExecutionSlot implements AutoCloseable { + + private final WorkerExecutionSlotPool owner; + private final AtomicBoolean closed = new AtomicBoolean(false); + + private WorkerExecutionSlot(WorkerExecutionSlotPool owner) { + this.owner = owner; + } + + @Override + public void close() { + if (closed.compareAndSet(false, true)) { + owner.release(); + } + } + } +} diff --git a/src/main/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerExecutionLifecycleManager.java b/src/main/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerExecutionLifecycleManager.java new file mode 100644 index 0000000..78cac15 --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerExecutionLifecycleManager.java @@ -0,0 +1,99 @@ +package com.opensource.docgrid.domain.worker.lifecycle; + +import static com.opensource.docgrid.domain.worker.config.WorkerExecutionConfig.WORKER_JOB_EXECUTOR; +import static com.opensource.docgrid.domain.worker.config.WorkerExecutionConfig.WORKER_LEASE_SCHEDULER; + +import java.time.Duration; +import java.util.List; +import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.context.event.ContextClosedEvent; +import org.springframework.context.event.EventListener; +import org.springframework.core.Ordered; +import org.springframework.core.annotation.Order; +import org.springframework.stereotype.Component; + +import com.opensource.docgrid.domain.worker.config.IndexingWorkerProperties; +import com.opensource.docgrid.domain.worker.service.WorkerLeaseRenewalManager; + +import lombok.extern.slf4j.Slf4j; + +/** + * Worker Context 종료 시 신규 Claim, 활성 실행과 Lease 갱신을 순서대로 정리한다. + * + *

Worker DB 상태가 STOPPED로 바뀌기 전에 Job Executor가 스스로 끝날 기회를 보장한다. 유예 시간을 + * 넘긴 실행은 Thread interrupt와 Lease 갱신 취소만 수행하며 DB 상태는 만료 복구 정책에 맡긴다. + */ +@Slf4j +@Component +@ConditionalOnProperty(prefix = "indexing.worker", name = "enabled", havingValue = "true") +public class WorkerExecutionLifecycleManager { + + private final WorkerJobPollingScheduler pollingScheduler; + private final ThreadPoolExecutor jobExecutor; + private final ScheduledThreadPoolExecutor leaseScheduler; + private final WorkerLeaseRenewalManager leaseRenewalManager; + private final IndexingWorkerProperties workerProperties; + private final AtomicBoolean closing = new AtomicBoolean(false); + + public WorkerExecutionLifecycleManager( + WorkerJobPollingScheduler pollingScheduler, + @Qualifier(WORKER_JOB_EXECUTOR) ThreadPoolExecutor jobExecutor, + @Qualifier(WORKER_LEASE_SCHEDULER) ScheduledThreadPoolExecutor leaseScheduler, + WorkerLeaseRenewalManager leaseRenewalManager, + IndexingWorkerProperties workerProperties + ) { + this.pollingScheduler = pollingScheduler; + this.jobExecutor = jobExecutor; + this.leaseScheduler = leaseScheduler; + this.leaseRenewalManager = leaseRenewalManager; + this.workerProperties = workerProperties; + } + + /** + * Worker STOPPED 기록보다 먼저 신규 Polling을 차단하고 활성 실행을 제한 시간 동안 기다린다. + */ + @Order(Ordered.HIGHEST_PRECEDENCE) + @EventListener(ContextClosedEvent.class) + public void shutdown() { + if (!closing.compareAndSet(false, true)) { + return; + } + + // 1. 신규 Slot과 Claim을 막은 뒤 이미 제출된 실행만 완료할 수 있게 Executor를 닫는다. + pollingScheduler.stopPolling(); + jobExecutor.shutdown(); + + // 2. 설정된 유예 시간 안에 모든 실행이 끝나면 interrupt 없이 Lease 자원만 정리한다. + boolean terminated = awaitJobTermination(workerProperties.getShutdownGracePeriod()); + List cancelledTasks = List.of(); + if (!terminated) { + // 3. 시간 초과 실행은 interrupt하고 DB 상태를 직접 변경하지 않아 Lease 복구가 회수하게 한다. + cancelledTasks = jobExecutor.shutdownNow(); + } + + // 4. Job Thread 정리 뒤 남은 갱신 Handle과 Scheduler를 닫고 Worker STOPPED Listener에 제어를 넘긴다. + leaseRenewalManager.stopAll(); + leaseScheduler.shutdown(); + log.info( + "Worker 실행 종료를 완료했습니다. graceful={}, cancelledTaskCount={}, activeLeaseCount={}", + terminated, + cancelledTasks.size(), + leaseRenewalManager.getActiveHandleCount() + ); + } + + private boolean awaitJobTermination(Duration gracePeriod) { + try { + return jobExecutor.awaitTermination(gracePeriod.toNanos(), TimeUnit.NANOSECONDS); + } catch (InterruptedException exception) { + Thread.currentThread().interrupt(); + return false; + } + } +} diff --git a/src/main/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerJobPollingScheduler.java b/src/main/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerJobPollingScheduler.java new file mode 100644 index 0000000..93ac1d5 --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerJobPollingScheduler.java @@ -0,0 +1,178 @@ +package com.opensource.docgrid.domain.worker.lifecycle; + +import static com.opensource.docgrid.domain.worker.config.WorkerExecutionConfig.WORKER_JOB_EXECUTOR; + +import java.util.Optional; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +import com.opensource.docgrid.domain.embedding.dto.response.ClaimedEmbeddingJobResponse; +import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobClaimService; +import com.opensource.docgrid.domain.worker.execution.WorkerExecutionSlotPool; +import com.opensource.docgrid.domain.worker.execution.WorkerExecutionSlotPool.WorkerExecutionSlot; +import com.opensource.docgrid.domain.worker.service.WorkerIndexingPipeline; +import com.opensource.docgrid.global.exception.DocGridException; + +import lombok.extern.slf4j.Slf4j; + +/** + * 등록된 Worker의 빈 실행 슬롯만큼 PENDING Job을 Claim해 제한 Executor에 제출한다. + * + *

실행 슬롯을 Claim 전에 확보하므로 로컬 처리 능력을 초과한 PROCESSING Job을 만들지 않는다. 빈 Claim, + * 조회 오류와 종료 경쟁은 현재 Polling 주기 안에서 격리해 Spring Scheduler의 다음 실행을 유지한다. + */ +@Slf4j +@Component +@ConditionalOnProperty(prefix = "indexing.worker", name = "enabled", havingValue = "true") +public class WorkerJobPollingScheduler { + + private final WorkerLifecycleManager workerLifecycleManager; + private final EmbeddingJobClaimService claimService; + private final WorkerExecutionSlotPool slotPool; + private final ThreadPoolExecutor jobExecutor; + private final WorkerIndexingPipeline indexingPipeline; + private final AtomicBoolean polling = new AtomicBoolean(false); + + public WorkerJobPollingScheduler( + WorkerLifecycleManager workerLifecycleManager, + EmbeddingJobClaimService claimService, + WorkerExecutionSlotPool slotPool, + @Qualifier(WORKER_JOB_EXECUTOR) ThreadPoolExecutor jobExecutor, + WorkerIndexingPipeline indexingPipeline + ) { + this.workerLifecycleManager = workerLifecycleManager; + this.claimService = claimService; + this.slotPool = slotPool; + this.jobExecutor = jobExecutor; + this.indexingPipeline = indexingPipeline; + } + + /** + * 현재 가용 슬롯 수만큼 Job을 찾아 즉시 실행하고 첫 빈 Claim에서 주기를 끝낸다. + */ + @Scheduled( + fixedDelayString = "${indexing.worker.polling-interval:1s}", + initialDelayString = "${indexing.worker.polling-interval:1s}" + ) + public void poll() { + if (!slotPool.isAccepting() || !polling.compareAndSet(false, true)) { + return; + } + + try { + Optional registeredWorkerId = workerLifecycleManager.getWorkerId(); + if (registeredWorkerId.isEmpty()) { + return; + } + + // 1. 한 주기에는 설정된 전체 슬롯 수까지만 Claim을 시도하고 사용 중 슬롯은 즉시 건너뛴다. + for (int index = 0; index < slotPool.getCapacity(); index++) { + Optional acquiredSlot = slotPool.tryAcquire(); + if (acquiredSlot.isEmpty()) { + return; + } + + // 2. 슬롯 획득 뒤 시작된 종료와 Claim 오류는 이 Slot만 반환하고 현재 주기를 끝낸다. + WorkerExecutionSlot executionSlot = acquiredSlot.get(); + if (!slotPool.isAccepting()) { + executionSlot.close(); + return; + } + Optional claimedJob = claim( + registeredWorkerId.get(), + executionSlot + ); + if (claimedJob.isEmpty()) { + return; + } + + // 3. Queue 없이 즉시 실행하며 종료 경쟁으로 제출이 거부되면 Slot만 반환한다. + if (!submit(claimedJob.get(), executionSlot)) { + return; + } + } + } finally { + polling.set(false); + } + } + + /** + * 종료 Listener가 신규 Slot과 Claim을 중단한다. + */ + public void stopPolling() { + slotPool.stopAccepting(); + } + + public boolean isPollingEnabled() { + return slotPool.isAccepting(); + } + + private Optional claim( + Long workerId, + WorkerExecutionSlot executionSlot + ) { + try { + Optional claimedJob = claimService.claim(workerId); + if (claimedJob.isEmpty()) { + executionSlot.close(); + } + return claimedJob; + } catch (RuntimeException exception) { + executionSlot.close(); + log.error( + "Worker Job Polling 중 Claim에 실패했습니다. workerId={}, errorCode={}", + workerId, + diagnosticCode(exception) + ); + return Optional.empty(); + } + } + + private boolean submit( + ClaimedEmbeddingJobResponse claimedJob, + WorkerExecutionSlot executionSlot + ) { + try { + jobExecutor.execute(() -> execute(claimedJob, executionSlot)); + return true; + } catch (RejectedExecutionException exception) { + executionSlot.close(); + log.warn( + "종료 중 Worker Job 실행 제출이 거부됐습니다. workerId={}, jobId={}", + claimedJob.workerId(), + claimedJob.jobId() + ); + return false; + } + } + + private void execute( + ClaimedEmbeddingJobResponse claimedJob, + WorkerExecutionSlot executionSlot + ) { + try { + indexingPipeline.execute(claimedJob, executionSlot); + } catch (RuntimeException exception) { + // Attempt 시작 이전 오류는 실패 API로 합성하지 않고 Lease 만료 복구가 현재 Claim을 회수하게 한다. + log.error( + "Worker Job 실행을 시작하지 못했습니다. workerId={}, jobId={}, errorCode={}", + claimedJob.workerId(), + claimedJob.jobId(), + diagnosticCode(exception) + ); + } + } + + private String diagnosticCode(RuntimeException exception) { + if (exception instanceof DocGridException docGridException) { + return docGridException.getErrorCode().getCode(); + } + return exception.getClass().getSimpleName(); + } +} diff --git a/src/main/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerLifecycleManager.java b/src/main/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerLifecycleManager.java index 4a394c8..03e4af6 100644 --- a/src/main/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerLifecycleManager.java +++ b/src/main/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerLifecycleManager.java @@ -10,6 +10,8 @@ import org.springframework.boot.context.event.ApplicationReadyEvent; import org.springframework.context.event.ContextClosedEvent; import org.springframework.context.event.EventListener; +import org.springframework.core.Ordered; +import org.springframework.core.annotation.Order; import org.springframework.stereotype.Component; import com.opensource.docgrid.domain.worker.config.IndexingWorkerProperties; @@ -18,6 +20,11 @@ import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; +/** + * 애플리케이션 실행 인스턴스를 Worker Node로 등록하고 Heartbeat·Job 실행에 현재 Worker ID를 제공한다. + * + *

종료 시 활성 Job Executor와 Lease 갱신이 먼저 정리된 뒤 Worker를 STOPPED로 기록한다. + */ @Slf4j @Component @RequiredArgsConstructor @@ -58,6 +65,7 @@ public void registerWorker() { ); } + @Order(Ordered.LOWEST_PRECEDENCE) @EventListener(ContextClosedEvent.class) public void stopWorker() { Long registeredWorkerId = workerId.getAndSet(null); diff --git a/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailure.java b/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailure.java new file mode 100644 index 0000000..a2a29ff --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailure.java @@ -0,0 +1,34 @@ +package com.opensource.docgrid.domain.worker.service; + +import com.opensource.docgrid.domain.embedding.enums.IndexingFailureType; + +/** + * Worker 실행 예외를 DB 실패 보고에 사용할 제한된 유형과 안전한 진단 정보로 표현한다. + * + *

소유권을 이미 잃은 예외는 {@code reportable=false}이며 실패 유형이 없다. 진단 코드와 메시지는 + * 자유 형식 원인이나 Claim Token을 포함하지 않고 로그와 실패 전이 경계에서만 사용한다. + */ +public record WorkerIndexingFailure( + boolean reportable, + IndexingFailureType failureType, + String diagnosticCode, + String safeMessage +) { + + static WorkerIndexingFailure reportable( + IndexingFailureType failureType, + String diagnosticCode, + String safeMessage + ) { + return new WorkerIndexingFailure(true, failureType, diagnosticCode, safeMessage); + } + + static WorkerIndexingFailure ownershipLost(String diagnosticCode) { + return new WorkerIndexingFailure( + false, + null, + diagnosticCode, + "현재 Worker 실행이 Job 소유권을 더 이상 보유하지 않습니다." + ); + } +} diff --git a/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailureClassifier.java b/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailureClassifier.java new file mode 100644 index 0000000..3945310 --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailureClassifier.java @@ -0,0 +1,111 @@ +package com.opensource.docgrid.domain.worker.service; + +import java.util.EnumSet; +import java.util.Set; + +import org.springframework.stereotype.Component; + +import com.opensource.docgrid.domain.embedding.enums.IndexingFailureType; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +/** + * Worker Pipeline 예외를 공개된 인덱싱 실패 정책과 안전한 고정 메시지로 분류한다. + * + *

예외 원문과 Stack Trace는 문서 내용이나 외부 연결 정보를 포함할 수 있으므로 DB 실패 메시지로 + * 전달하지 않는다. 소유권 상실 오류는 과거 Claim으로 상태를 덮어쓰지 않도록 보고 대상에서 제외한다. + */ +@Component +public class WorkerIndexingFailureClassifier { + + private static final Set OWNERSHIP_LOST_ERRORS = EnumSet.of( + ErrorCode.EMBEDDING_JOB_NOT_FOUND, + ErrorCode.EMBEDDING_JOB_NOT_PROCESSING, + ErrorCode.EMBEDDING_JOB_OWNERSHIP_INVALID, + ErrorCode.EMBEDDING_JOB_LEASE_EXPIRED, + ErrorCode.EMBEDDING_JOB_ATTEMPT_INVALID, + ErrorCode.EMBEDDING_JOB_FAILURE_CONFLICT + ); + private static final Set DOCUMENT_CONTENT_ERRORS = EnumSet.of( + ErrorCode.UNSUPPORTED_DOCUMENT_TYPE, + ErrorCode.DOCUMENT_CONTENT_EMPTY, + ErrorCode.DOCUMENT_TEXT_DECODING_FAILED + ); + private static final Set EMBEDDING_RESULT_ERRORS = EnumSet.of( + ErrorCode.EMBEDDING_DIMENSION_MISMATCH, + ErrorCode.EMBEDDING_VECTOR_INVALID + ); + private static final Set INDEXING_STATE_ERRORS = EnumSet.of( + ErrorCode.DOCUMENT_VERSION_CHUNKING_NOT_ALLOWED, + ErrorCode.DOCUMENT_VERSION_EMBEDDING_NOT_ALLOWED, + ErrorCode.INDEXING_STATUS_INCONSISTENT, + ErrorCode.DOCUMENT_FILE_REFERENCE_MISSING, + ErrorCode.DOCUMENT_CHUNKS_INCONSISTENT, + ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT, + ErrorCode.DOCUMENT_INDEXING_COMPLETION_NOT_ALLOWED, + ErrorCode.DOCUMENT_INDEXING_STALE_COMPLETION, + ErrorCode.DOCUMENT_INDEXING_COMPLETION_INCONSISTENT, + ErrorCode.DOCUMENT_INDEXING_FAILURE_INCONSISTENT, + ErrorCode.EMBEDDING_JOB_OWNERSHIP_INCONSISTENT, + ErrorCode.EMBEDDING_MODEL_NOT_CONFIGURED, + ErrorCode.MULTIPLE_ACTIVE_EMBEDDING_MODELS + ); + + /** + * 실행 예외를 실패 보고 가능 여부, 제한 유형과 고정 진단 메시지로 변환한다. + */ + public WorkerIndexingFailure classify(RuntimeException exception) { + if (!(exception instanceof DocGridException docGridException)) { + return WorkerIndexingFailure.reportable( + IndexingFailureType.WORKER_INTERNAL_ERROR, + "UNEXPECTED_RUNTIME_EXCEPTION", + "Worker 내부 실행 오류로 인덱싱을 완료하지 못했습니다." + ); + } + + ErrorCode errorCode = docGridException.getErrorCode(); + if (OWNERSHIP_LOST_ERRORS.contains(errorCode)) { + return WorkerIndexingFailure.ownershipLost(errorCode.getCode()); + } + if (errorCode == ErrorCode.FILE_STORAGE_FAILED) { + return WorkerIndexingFailure.reportable( + IndexingFailureType.STORAGE_UNAVAILABLE, + errorCode.getCode(), + "파일 저장소를 사용할 수 없어 인덱싱을 완료하지 못했습니다." + ); + } + if (DOCUMENT_CONTENT_ERRORS.contains(errorCode)) { + return WorkerIndexingFailure.reportable( + IndexingFailureType.DOCUMENT_CONTENT_INVALID, + errorCode.getCode(), + "문서 내용을 인덱싱 가능한 텍스트로 처리할 수 없습니다." + ); + } + if (errorCode == ErrorCode.EMBEDDING_SERVER_UNAVAILABLE) { + return WorkerIndexingFailure.reportable( + IndexingFailureType.EMBEDDING_PROVIDER_UNAVAILABLE, + errorCode.getCode(), + "Embedding Provider를 사용할 수 없어 인덱싱을 완료하지 못했습니다." + ); + } + if (EMBEDDING_RESULT_ERRORS.contains(errorCode)) { + return WorkerIndexingFailure.reportable( + IndexingFailureType.EMBEDDING_RESULT_INVALID, + errorCode.getCode(), + "Embedding Provider 결과가 현재 모델 계약과 일치하지 않습니다." + ); + } + if (INDEXING_STATE_ERRORS.contains(errorCode)) { + return WorkerIndexingFailure.reportable( + IndexingFailureType.INDEXING_STATE_INCONSISTENT, + errorCode.getCode(), + "인덱싱 Job과 문서 파이프라인 상태가 일치하지 않습니다." + ); + } + return WorkerIndexingFailure.reportable( + IndexingFailureType.WORKER_INTERNAL_ERROR, + errorCode.getCode(), + "Worker 내부 실행 오류로 인덱싱을 완료하지 못했습니다." + ); + } +} diff --git a/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailureReporter.java b/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailureReporter.java new file mode 100644 index 0000000..c08cc48 --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailureReporter.java @@ -0,0 +1,85 @@ +package com.opensource.docgrid.domain.worker.service; + +import org.springframework.stereotype.Service; + +import com.opensource.docgrid.domain.embedding.dto.request.FailDocumentIndexingRequest; +import com.opensource.docgrid.domain.embedding.dto.response.ClaimedEmbeddingJobResponse; +import com.opensource.docgrid.domain.embedding.service.command.DocumentIndexingFailureService; +import com.opensource.docgrid.global.exception.DocGridException; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; + +/** + * Attempt 시작 이후 Worker Pipeline 오류를 기존 인덱싱 실패 전이 Service에 안전하게 보고한다. + * + *

분류된 고정 메시지만 저장하고 소유권 상실 오류는 보고하지 않는다. 보고 자체의 실패도 비동기 Job + * Thread 밖으로 전파하지 않아 남은 상태는 Lease 만료 복구가 처리할 수 있게 한다. + */ +@Slf4j +@Service +@RequiredArgsConstructor +public class WorkerIndexingFailureReporter { + + private final WorkerIndexingFailureClassifier failureClassifier; + private final DocumentIndexingFailureService failureService; + + /** + * 현재 Attempt 오류를 분류해 보고하고, 과거 소유권이거나 보고가 실패하면 안전한 진단만 남긴다. + */ + public void report( + ClaimedEmbeddingJobResponse claimedJob, + Long attemptId, + RuntimeException exception + ) { + WorkerIndexingFailure failure = failureClassifier.classify(exception); + if (!failure.reportable()) { + log.warn( + "소유권을 잃은 Worker 인덱싱 실행을 중단합니다. workerId={}, jobId={}, attemptId={}, errorCode={}", + claimedJob.workerId(), + claimedJob.jobId(), + attemptId, + failure.diagnosticCode() + ); + return; + } + + try { + failureService.fail( + claimedJob.jobId(), + attemptId, + new FailDocumentIndexingRequest( + claimedJob.workerId(), + claimedJob.claimToken(), + failure.failureType(), + failure.safeMessage() + ) + ); + log.info( + "Worker 인덱싱 실패를 기록했습니다. workerId={}, jobId={}, attemptId={}, failureType={}, errorCode={}", + claimedJob.workerId(), + claimedJob.jobId(), + attemptId, + failure.failureType(), + failure.diagnosticCode() + ); + } catch (RuntimeException reportException) { + log.error( + "Worker 인덱싱 실패를 기록하지 못했습니다. workerId={}, jobId={}, attemptId={}, " + + "originalErrorCode={}, reportErrorCode={}", + claimedJob.workerId(), + claimedJob.jobId(), + attemptId, + failure.diagnosticCode(), + diagnosticCode(reportException) + ); + } + } + + private String diagnosticCode(RuntimeException exception) { + if (exception instanceof DocGridException docGridException) { + return docGridException.getErrorCode().getCode(); + } + return exception.getClass().getSimpleName(); + } +} diff --git a/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingPipeline.java b/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingPipeline.java new file mode 100644 index 0000000..dbadd3e --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingPipeline.java @@ -0,0 +1,141 @@ +package com.opensource.docgrid.domain.worker.service; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Service; + +import com.opensource.docgrid.domain.document.enums.DocumentVersionStatus; +import com.opensource.docgrid.domain.document.service.DocumentParsingService; +import com.opensource.docgrid.domain.document.service.query.DocumentIndexingStageQueryService; +import com.opensource.docgrid.domain.embedding.dto.request.CompleteDocumentIndexingRequest; +import com.opensource.docgrid.domain.embedding.dto.request.CreateDocumentChunksRequest; +import com.opensource.docgrid.domain.embedding.dto.request.CreateDocumentEmbeddingsRequest; +import com.opensource.docgrid.domain.embedding.dto.request.StartEmbeddingJobAttemptRequest; +import com.opensource.docgrid.domain.embedding.dto.response.ClaimedEmbeddingJobResponse; +import com.opensource.docgrid.domain.embedding.dto.response.StartedEmbeddingJobAttemptResponse; +import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus; +import com.opensource.docgrid.domain.embedding.service.DocumentEmbeddingService; +import com.opensource.docgrid.domain.embedding.service.command.DocumentIndexingCompletionService; +import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobAttemptService; +import com.opensource.docgrid.domain.worker.execution.WorkerExecutionSlotPool.WorkerExecutionSlot; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; + +/** + * Claim된 Embedding Job을 기존 인덱싱 Service들로 끝까지 실행하는 Worker 오케스트레이터다. + * + *

자체 Transaction을 만들지 않아 파일 저장소와 Embedding Provider 호출 중 DB 잠금을 유지하지 않는다. + * 문서 버전 상태 Snapshot은 시작 단계만 결정하고 각 Service가 현재 Claim과 상태를 다시 검증한다. + */ +@Slf4j +@Service +@RequiredArgsConstructor +@ConditionalOnProperty(prefix = "indexing.worker", name = "enabled", havingValue = "true") +public class WorkerIndexingPipeline { + + private final EmbeddingJobAttemptService embeddingJobAttemptService; + private final DocumentIndexingStageQueryService stageQueryService; + private final DocumentParsingService documentParsingService; + private final DocumentEmbeddingService documentEmbeddingService; + private final DocumentIndexingCompletionService completionService; + private final WorkerIndexingFailureReporter failureReporter; + private final WorkerLeaseRenewalManager leaseRenewalManager; + + /** + * 현재 Claim의 Attempt를 시작하고 문서 상태에 맞는 단계부터 완료까지 실행한다. + * + *

호출자가 넘긴 실행 Slot의 소유권을 인수하며 성공과 예외 경로 모두에서 정확히 한 번 반환한다. + */ + public void execute( + ClaimedEmbeddingJobResponse claimedJob, + WorkerExecutionSlot executionSlot + ) { + try (executionSlot) { + validateClaim(claimedJob); + + // 1. 실제 실행 Context를 먼저 기록해 이후 단계와 실패 보고가 같은 Attempt를 식별하게 한다. + StartedEmbeddingJobAttemptResponse attempt = embeddingJobAttemptService.start( + claimedJob.jobId(), + new StartEmbeddingJobAttemptRequest(claimedJob.workerId(), claimedJob.claimToken()) + ).response(); + + // 2. 실제 Attempt 수명에 맞춰 Lease 갱신을 시작하고 모든 종료 경로에서 예약을 해제한다. + try (WorkerLeaseRenewalHandle leaseHandle = leaseRenewalManager.start(claimedJob)) { + leaseHandle.ensureOwned(); + + // 3. 경로 선택용 상태 Snapshot을 조회하고 이미 완료한 단계는 다시 외부 호출하지 않는다. + DocumentVersionStatus initialStatus = stageQueryService.getStatus( + claimedJob.documentVersionId() + ); + log.info( + "Worker 인덱싱 실행을 시작합니다. workerId={}, jobId={}, attemptId={}, initialStatus={}", + claimedJob.workerId(), + claimedJob.jobId(), + attempt.attemptId(), + initialStatus + ); + executeFromCurrentStage(claimedJob, attempt.attemptId(), initialStatus, leaseHandle); + + // 4. 전체 Embedding Set과 현재 실행 소유권을 최종 검증해 검색 가능한 Version으로 확정한다. + leaseHandle.ensureOwned(); + completionService.complete( + claimedJob.jobId(), + attempt.attemptId(), + new CompleteDocumentIndexingRequest(claimedJob.workerId(), claimedJob.claimToken()) + ); + log.info( + "Worker 인덱싱 실행을 완료했습니다. workerId={}, jobId={}, attemptId={}", + claimedJob.workerId(), + claimedJob.jobId(), + attempt.attemptId() + ); + } catch (RuntimeException exception) { + // 5. 실제 Attempt가 시작된 뒤의 오류만 제한된 실패 계약으로 기록한다. + failureReporter.report(claimedJob, attempt.attemptId(), exception); + } + } + } + + private void validateClaim(ClaimedEmbeddingJobResponse claimedJob) { + if (claimedJob == null + || claimedJob.status() != EmbeddingJobStatus.PROCESSING + || claimedJob.jobId() == null + || claimedJob.workerId() == null + || claimedJob.documentVersionId() == null + || claimedJob.claimToken() == null) { + throw new DocGridException(ErrorCode.EMBEDDING_JOB_OWNERSHIP_INCONSISTENT); + } + } + + private void executeFromCurrentStage( + ClaimedEmbeddingJobResponse claimedJob, + Long attemptId, + DocumentVersionStatus initialStatus, + WorkerLeaseRenewalHandle leaseHandle + ) { + // 1. 새 Version과 중단된 Parsing은 기존 Chunk Service의 멱등·재개 계약으로 CHUNKED까지 진행한다. + if (initialStatus == DocumentVersionStatus.UPLOADED + || initialStatus == DocumentVersionStatus.PARSING) { + documentParsingService.createChunks( + claimedJob.jobId(), + attemptId, + new CreateDocumentChunksRequest(claimedJob.workerId(), claimedJob.claimToken()) + ); + leaseHandle.ensureOwned(); + } else if (initialStatus != DocumentVersionStatus.CHUNKED + && initialStatus != DocumentVersionStatus.EMBEDDING) { + throw new DocGridException(ErrorCode.INDEXING_STATUS_INCONSISTENT); + } + + // 2. CHUNKED 또는 중단된 EMBEDDING 상태를 기존 Service의 생성·재생 계약으로 완료한다. + leaseHandle.ensureOwned(); + documentEmbeddingService.createEmbeddings( + claimedJob.jobId(), + attemptId, + new CreateDocumentEmbeddingsRequest(claimedJob.workerId(), claimedJob.claimToken()) + ); + leaseHandle.ensureOwned(); + } +} diff --git a/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerLeaseRenewalHandle.java b/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerLeaseRenewalHandle.java new file mode 100644 index 0000000..2c33eea --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerLeaseRenewalHandle.java @@ -0,0 +1,102 @@ +package com.opensource.docgrid.domain.worker.service; + +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Consumer; + +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +/** + * 한 Worker Job 실행의 Lease 갱신 예약과 소유권 상태를 보관하는 수명 Handle이다. + * + *

Pipeline만 이 Handle을 닫으며 Scheduler는 권위 있는 갱신 거부가 발생하면 lost 상태로 전환한다. + * Claim Token은 외부에 노출하지 않고 갱신 Task 내부 요청 생성에만 사용한다. + */ +public final class WorkerLeaseRenewalHandle implements AutoCloseable { + + private final Long jobId; + private final Long workerId; + private final String claimToken; + private final Consumer inactiveCallback; + private final AtomicReference> scheduledFuture = new AtomicReference<>(); + private final AtomicBoolean ownershipLost = new AtomicBoolean(false); + private final AtomicBoolean closed = new AtomicBoolean(false); + + WorkerLeaseRenewalHandle( + Long jobId, + Long workerId, + String claimToken, + Consumer inactiveCallback + ) { + this.jobId = jobId; + this.workerId = workerId; + this.claimToken = claimToken; + this.inactiveCallback = inactiveCallback; + } + + /** + * Manager가 만든 주기 예약을 한 번만 연결한다. + */ + void attach(ScheduledFuture future) { + if (!scheduledFuture.compareAndSet(null, future)) { + future.cancel(false); + throw new IllegalStateException("Lease 갱신 예약은 한 번만 연결할 수 있습니다."); + } + if (closed.get() || ownershipLost.get()) { + future.cancel(false); + } + } + + /** + * 단계 시작 전에 Scheduler가 확인한 소유권 상실을 Pipeline에 전달한다. + */ + public void ensureOwned() { + if (ownershipLost.get() || closed.get()) { + throw new DocGridException(ErrorCode.EMBEDDING_JOB_OWNERSHIP_INVALID); + } + } + + void markOwnershipLost() { + if (ownershipLost.compareAndSet(false, true)) { + cancelScheduledTask(); + inactiveCallback.accept(this); + } + } + + boolean isClosed() { + return closed.get(); + } + + boolean isOwnershipLost() { + return ownershipLost.get(); + } + + Long jobId() { + return jobId; + } + + Long workerId() { + return workerId; + } + + String claimToken() { + return claimToken; + } + + @Override + public void close() { + if (closed.compareAndSet(false, true)) { + cancelScheduledTask(); + inactiveCallback.accept(this); + } + } + + private void cancelScheduledTask() { + ScheduledFuture future = scheduledFuture.get(); + if (future != null) { + future.cancel(false); + } + } +} diff --git a/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerLeaseRenewalManager.java b/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerLeaseRenewalManager.java new file mode 100644 index 0000000..a7b3819 --- /dev/null +++ b/src/main/java/com/opensource/docgrid/domain/worker/service/WorkerLeaseRenewalManager.java @@ -0,0 +1,157 @@ +package com.opensource.docgrid.domain.worker.service; + +import static com.opensource.docgrid.domain.worker.config.WorkerExecutionConfig.WORKER_LEASE_SCHEDULER; + +import java.time.Duration; +import java.util.EnumSet; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.stereotype.Service; + +import com.opensource.docgrid.domain.embedding.dto.request.RenewEmbeddingJobLeaseRequest; +import com.opensource.docgrid.domain.embedding.dto.response.ClaimedEmbeddingJobResponse; +import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobLeaseService; +import com.opensource.docgrid.domain.worker.config.IndexingWorkerProperties; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +import lombok.extern.slf4j.Slf4j; + +/** + * 활성 Worker Job별 Lease 갱신 Task를 등록하고 소유권 상실과 종료를 조정한다. + * + *

하나의 Scheduled Executor에서 짧은 갱신 Transaction만 실행한다. 권위 있는 소유권 오류는 해당 + * Handle을 lost로 전환하고, 일시적인 RuntimeException은 다음 갱신 주기와 단계별 DB fencing에 맡긴다. + */ +@Slf4j +@Service +@ConditionalOnProperty(prefix = "indexing.worker", name = "enabled", havingValue = "true") +public class WorkerLeaseRenewalManager { + + private static final Set OWNERSHIP_LOST_ERRORS = EnumSet.of( + ErrorCode.WORKER_NOT_FOUND, + ErrorCode.WORKER_NOT_AVAILABLE, + ErrorCode.EMBEDDING_JOB_NOT_FOUND, + ErrorCode.EMBEDDING_JOB_NOT_PROCESSING, + ErrorCode.EMBEDDING_JOB_OWNERSHIP_INVALID, + ErrorCode.EMBEDDING_JOB_LEASE_EXPIRED, + ErrorCode.EMBEDDING_JOB_OWNERSHIP_INCONSISTENT + ); + + private final ScheduledThreadPoolExecutor leaseScheduler; + private final EmbeddingJobLeaseService leaseService; + private final IndexingWorkerProperties workerProperties; + private final ConcurrentHashMap activeHandles = + new ConcurrentHashMap<>(); + private final AtomicBoolean accepting = new AtomicBoolean(true); + + public WorkerLeaseRenewalManager( + @Qualifier(WORKER_LEASE_SCHEDULER) ScheduledThreadPoolExecutor leaseScheduler, + EmbeddingJobLeaseService leaseService, + IndexingWorkerProperties workerProperties + ) { + this.leaseScheduler = leaseScheduler; + this.leaseService = leaseService; + this.workerProperties = workerProperties; + } + + /** + * Attempt가 시작된 Claim의 Lease를 설정 주기로 갱신하는 Handle을 등록한다. + */ + public WorkerLeaseRenewalHandle start(ClaimedEmbeddingJobResponse claimedJob) { + if (!accepting.get()) { + throw new IllegalStateException("Worker Lease 갱신 관리자가 종료 중입니다."); + } + + WorkerLeaseRenewalHandle handle = new WorkerLeaseRenewalHandle( + claimedJob.jobId(), + claimedJob.workerId(), + claimedJob.claimToken(), + inactiveHandle -> activeHandles.remove(claimedJob.jobId(), inactiveHandle) + ); + WorkerLeaseRenewalHandle existing = activeHandles.putIfAbsent(claimedJob.jobId(), handle); + if (existing != null) { + throw new DocGridException(ErrorCode.EMBEDDING_JOB_OWNERSHIP_INCONSISTENT); + } + if (!accepting.get()) { + handle.close(); + throw new IllegalStateException("Worker Lease 갱신 관리자가 종료 중입니다."); + } + + try { + Duration interval = workerProperties.getLeaseRenewalInterval(); + long intervalNanos = interval.toNanos(); + ScheduledFuture future = leaseScheduler.scheduleWithFixedDelay( + () -> renew(handle), + intervalNanos, + intervalNanos, + TimeUnit.NANOSECONDS + ); + handle.attach(future); + return handle; + } catch (RuntimeException exception) { + handle.close(); + throw exception; + } + } + + /** + * 애플리케이션 강제 종료 시 남은 모든 실행의 갱신 예약을 취소한다. + */ + public void stopAll() { + accepting.set(false); + activeHandles.values().forEach(WorkerLeaseRenewalHandle::close); + } + + public int getActiveHandleCount() { + return activeHandles.size(); + } + + public boolean isAccepting() { + return accepting.get(); + } + + private void renew(WorkerLeaseRenewalHandle handle) { + if (handle.isClosed() || handle.isOwnershipLost()) { + return; + } + + try { + leaseService.renew( + handle.jobId(), + new RenewEmbeddingJobLeaseRequest(handle.workerId(), handle.claimToken()) + ); + } catch (DocGridException exception) { + if (OWNERSHIP_LOST_ERRORS.contains(exception.getErrorCode())) { + handle.markOwnershipLost(); + log.warn( + "활성 Worker Job Lease 소유권을 잃었습니다. workerId={}, jobId={}, errorCode={}", + handle.workerId(), + handle.jobId(), + exception.getErrorCode().getCode() + ); + return; + } + log.error( + "활성 Worker Job Lease를 갱신하지 못했습니다. workerId={}, jobId={}, errorCode={}", + handle.workerId(), + handle.jobId(), + exception.getErrorCode().getCode() + ); + } catch (RuntimeException exception) { + log.error( + "활성 Worker Job Lease 갱신 중 Runtime 오류가 발생했습니다. workerId={}, jobId={}, errorType={}", + handle.workerId(), + handle.jobId(), + exception.getClass().getSimpleName() + ); + } + } +} diff --git a/src/main/resources/application.yml b/src/main/resources/application.yml index ea5234a..4024026 100644 --- a/src/main/resources/application.yml +++ b/src/main/resources/application.yml @@ -9,6 +9,11 @@ spring: multipart: max-file-size: 10MB max-request-size: 11MB + task: + scheduling: + pool: + # DB Claim 지연이 Worker Heartbeat와 Lease 복구 실행을 막지 않도록 세 작업을 분리한다. + size: 3 document: upload: @@ -23,12 +28,16 @@ indexing: name: ${INDEXING_WORKER_NAME:indexing-worker} heartbeat-interval: ${INDEXING_WORKER_HEARTBEAT_INTERVAL:10s} dead-threshold: ${INDEXING_WORKER_DEAD_THRESHOLD:30s} + polling-interval: ${INDEXING_WORKER_POLLING_INTERVAL:1s} + max-concurrency: ${INDEXING_WORKER_MAX_CONCURRENCY:2} lease-duration: ${INDEXING_WORKER_LEASE_DURATION:5m} + lease-renewal-interval: ${INDEXING_WORKER_LEASE_RENEWAL_INTERVAL:1m} lease-recovery-interval: ${INDEXING_WORKER_LEASE_RECOVERY_INTERVAL:30s} lease-recovery-batch-size: ${INDEXING_WORKER_LEASE_RECOVERY_BATCH_SIZE:100} # Retry 지연은 10초부터 지수 증가하며 운영 기본 상한인 5분에서 제한한다. retry-initial-delay: ${INDEXING_WORKER_RETRY_INITIAL_DELAY:10s} retry-max-delay: ${INDEXING_WORKER_RETRY_MAX_DELAY:5m} + shutdown-grace-period: ${INDEXING_WORKER_SHUTDOWN_GRACE_PERIOD:30s} server: port: 8080 diff --git a/src/test/java/com/opensource/docgrid/domain/worker/config/IndexingWorkerPropertiesTest.java b/src/test/java/com/opensource/docgrid/domain/worker/config/IndexingWorkerPropertiesTest.java index 93e4d23..638cfac 100644 --- a/src/test/java/com/opensource/docgrid/domain/worker/config/IndexingWorkerPropertiesTest.java +++ b/src/test/java/com/opensource/docgrid/domain/worker/config/IndexingWorkerPropertiesTest.java @@ -13,7 +13,7 @@ import jakarta.validation.Validator; /** - * Worker Lease와 Retry 지연 설정의 기본값 및 애플리케이션 시작 단계 유효성 검사를 검증하는 단위 테스트. + * Worker Polling, 실행 동시성, Lease와 종료 설정의 기본값 및 시작 단계 유효성 검사를 검증한다. * *

정상적인 양수 기간은 허용하고 발급 즉시 만료되는 0 또는 음수 기간은 차단하는지 확인한다. */ @@ -22,6 +22,50 @@ class IndexingWorkerPropertiesTest { private final Validator validator = Validation.buildDefaultValidatorFactory().getValidator(); + @Test + @DisplayName("기본 Polling과 실행 설정은 1초, 동시 실행 2, 종료 유예 30초다") + void defaultExecutionSettings_areValid() { + IndexingWorkerProperties properties = new IndexingWorkerProperties(); + + assertThat(properties.getPollingInterval()).isEqualTo(Duration.ofSeconds(1)); + assertThat(properties.getMaxConcurrency()).isEqualTo(2); + assertThat(properties.getShutdownGracePeriod()).isEqualTo(Duration.ofSeconds(30)); + assertThat(properties.isPollingIntervalValid()).isTrue(); + assertThat(properties.isShutdownGracePeriodValid()).isTrue(); + } + + @Test + @DisplayName("Polling 주기는 양수이고 종료 유예 시간은 0 이상이어야 한다") + void executionIntervals_areInvalid_when_outOfRange() { + IndexingWorkerProperties properties = new IndexingWorkerProperties(); + + properties.setPollingInterval(Duration.ZERO); + assertThat(properties.isPollingIntervalValid()).isFalse(); + + properties.setPollingInterval(Duration.ofMillis(1)); + properties.setShutdownGracePeriod(Duration.ZERO); + assertThat(properties.isShutdownGracePeriodValid()).isTrue(); + + properties.setShutdownGracePeriod(Duration.ofNanos(-1)); + assertThat(properties.isShutdownGracePeriodValid()).isFalse(); + } + + @Test + @DisplayName("Lease 갱신 주기는 양수이고 Lease 기간보다 짧아야 한다") + void leaseRenewalInterval_isValid_onlyBeforeLeaseExpiry() { + IndexingWorkerProperties properties = new IndexingWorkerProperties(); + properties.setLeaseDuration(Duration.ofMinutes(5)); + + properties.setLeaseRenewalInterval(Duration.ofMinutes(1)); + assertThat(properties.isLeaseRenewalIntervalValid()).isTrue(); + + properties.setLeaseRenewalInterval(Duration.ZERO); + assertThat(properties.isLeaseRenewalIntervalValid()).isFalse(); + + properties.setLeaseRenewalInterval(Duration.ofMinutes(5)); + assertThat(properties.isLeaseRenewalIntervalValid()).isFalse(); + } + @Test @DisplayName("기본 Lease 기간은 5분이며 유효하다") void defaultLeaseDuration_isValid() { diff --git a/src/test/java/com/opensource/docgrid/domain/worker/execution/WorkerExecutionSlotPoolTest.java b/src/test/java/com/opensource/docgrid/domain/worker/execution/WorkerExecutionSlotPoolTest.java new file mode 100644 index 0000000..5c9f7a4 --- /dev/null +++ b/src/test/java/com/opensource/docgrid/domain/worker/execution/WorkerExecutionSlotPoolTest.java @@ -0,0 +1,63 @@ +package com.opensource.docgrid.domain.worker.execution; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import com.opensource.docgrid.domain.worker.execution.WorkerExecutionSlotPool.WorkerExecutionSlot; + +/** + * Worker 로컬 실행 슬롯의 최대 동시성, 단일 반환과 종료 후 획득 차단을 검증한다. + */ +@DisplayName("WorkerExecutionSlotPool 테스트") +class WorkerExecutionSlotPoolTest { + + @Test + @DisplayName("설정된 수까지만 Slot을 획득하고 close 시 다시 사용할 수 있다") + void tryAcquire_limitsConcurrencyAndReleasesSlot() { + WorkerExecutionSlotPool slotPool = new WorkerExecutionSlotPool(2); + WorkerExecutionSlot first = slotPool.tryAcquire().orElseThrow(); + WorkerExecutionSlot second = slotPool.tryAcquire().orElseThrow(); + + assertThat(slotPool.getActiveSlots()).isEqualTo(2); + assertThat(slotPool.tryAcquire()).isEmpty(); + + first.close(); + assertThat(slotPool.getAvailableSlots()).isOne(); + assertThat(slotPool.tryAcquire()).isPresent(); + second.close(); + } + + @Test + @DisplayName("같은 Slot을 여러 번 닫아도 Permit은 한 번만 반환한다") + void close_releasesPermitOnlyOnce() { + WorkerExecutionSlotPool slotPool = new WorkerExecutionSlotPool(1); + WorkerExecutionSlot slot = slotPool.tryAcquire().orElseThrow(); + + slot.close(); + slot.close(); + + assertThat(slotPool.getAvailableSlots()).isOne(); + assertThat(slotPool.getActiveSlots()).isZero(); + } + + @Test + @DisplayName("Polling 중단 뒤에는 남은 Permit이 있어도 신규 Slot을 주지 않는다") + void stopAccepting_blocksNewSlots() { + WorkerExecutionSlotPool slotPool = new WorkerExecutionSlotPool(1); + + slotPool.stopAccepting(); + + assertThat(slotPool.isAccepting()).isFalse(); + assertThat(slotPool.tryAcquire()).isEmpty(); + } + + @Test + @DisplayName("실행 슬롯 수는 1 이상이어야 한다") + void constructor_rejectsNonPositiveCapacity() { + assertThatThrownBy(() -> new WorkerExecutionSlotPool(0)) + .isInstanceOf(IllegalArgumentException.class); + } +} diff --git a/src/test/java/com/opensource/docgrid/domain/worker/integration/WorkerOrchestrationIntegrationTest.java b/src/test/java/com/opensource/docgrid/domain/worker/integration/WorkerOrchestrationIntegrationTest.java new file mode 100644 index 0000000..12e10fc --- /dev/null +++ b/src/test/java/com/opensource/docgrid/domain/worker/integration/WorkerOrchestrationIntegrationTest.java @@ -0,0 +1,524 @@ +package com.opensource.docgrid.domain.worker.integration; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.then; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; + +import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.List; +import java.util.Optional; +import java.util.UUID; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.SynchronousQueue; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; + +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestInstance; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.ActiveProfiles; +import org.springframework.test.context.DynamicPropertyRegistry; +import org.springframework.test.context.DynamicPropertySource; + +import com.opensource.docgrid.domain.embedding.dto.response.ClaimedEmbeddingJobResponse; +import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobClaimService; +import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobLeaseRecoveryService; +import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobLeaseRecoveryService.RecoveryResult; +import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobLeaseService; +import com.opensource.docgrid.domain.worker.config.IndexingWorkerProperties; +import com.opensource.docgrid.domain.worker.execution.WorkerExecutionSlotPool; +import com.opensource.docgrid.domain.worker.execution.WorkerExecutionSlotPool.WorkerExecutionSlot; +import com.opensource.docgrid.domain.worker.lifecycle.WorkerJobPollingScheduler; +import com.opensource.docgrid.domain.worker.lifecycle.WorkerLifecycleManager; +import com.opensource.docgrid.domain.worker.service.WorkerIndexingPipeline; +import com.opensource.docgrid.domain.worker.service.WorkerLeaseRenewalHandle; +import com.opensource.docgrid.domain.worker.service.WorkerLeaseRenewalManager; + +/** + * 실제 PostgreSQL에서 Worker Polling 슬롯과 Lease 갱신·복구의 교차 계층 동시성 경계를 검증한다. + * + *

Poller는 실제 Claim Service를 사용하되 외부 I/O Pipeline은 제어 가능한 Mock으로 대체한다. Lease + * 시나리오는 실제 갱신과 복구 Transaction을 사용해 한 Job의 DB 소유권이 중복 처리되지 않는지 확인한다. + */ +@Tag("integration") +@Tag("claim-concurrency") +@ActiveProfiles("test") +@SpringBootTest +@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS) +@TestInstance(TestInstance.Lifecycle.PER_CLASS) +@DisplayName("Worker 오케스트레이션 PostgreSQL 통합 테스트") +class WorkerOrchestrationIntegrationTest { + + private static final String TEST_SCHEMA = "docgrid_worker_orchestration_integration_test"; + private static final int MAX_CONCURRENCY = 2; + private static final long TIMEOUT_SECONDS = 10; + + @Autowired private JdbcTemplate jdbcTemplate; + @Autowired private EmbeddingJobClaimService claimService; + @Autowired private EmbeddingJobLeaseService leaseService; + @Autowired private EmbeddingJobLeaseRecoveryService recoveryService; + @Autowired private IndexingWorkerProperties workerProperties; + + @DynamicPropertySource + static void configureDatabase(DynamicPropertyRegistry registry) { + registry.add("TEST_DB_SCHEMA", () -> TEST_SCHEMA); + registry.add("jwt.secret", () -> "docgrid-worker-orchestration-integration-test-secret-key-2026"); + registry.add("indexing.worker.dead-threshold", () -> "10m"); + registry.add("indexing.worker.lease-duration", () -> "2s"); + registry.add("indexing.worker.lease-renewal-interval", () -> "50ms"); + registry.add("indexing.worker.retry-initial-delay", () -> "10s"); + registry.add("indexing.worker.retry-max-delay", () -> "5m"); + } + + @BeforeEach + void resetState() { + jdbcTemplate.execute(""" + TRUNCATE TABLE + indexing_events, + embedding_job_attempts, + embedding_jobs, + document_versions, + documents, + worker_nodes, + users + RESTART IDENTITY CASCADE + """); + } + + @AfterAll + void dropSchema() { + jdbcTemplate.execute("DROP SCHEMA IF EXISTS " + TEST_SCHEMA + " CASCADE"); + } + + @Test + @DisplayName("여러 Poller가 한 Job을 경쟁해도 실제 Claim과 실행 제출은 한 번만 발생한다") + void poll_concurrentWorkersClaimOneJobExactlyOnce() throws Exception { + TestContext context = insertContext("single-claim"); + Long jobId = insertPendingJobs(context, 1).get(0); + Long firstWorkerId = insertActiveWorker("poller-a"); + Long secondWorkerId = insertActiveWorker("poller-b"); + WorkerIndexingPipeline pipeline = mock(WorkerIndexingPipeline.class); + CountDownLatch executed = new CountDownLatch(1); + doAnswer(invocation -> { + WorkerExecutionSlot slot = invocation.getArgument(1); + try { + executed.countDown(); + } finally { + slot.close(); + } + return null; + }).when(pipeline).execute(any(ClaimedEmbeddingJobResponse.class), any(WorkerExecutionSlot.class)); + + ThreadPoolExecutor firstJobExecutor = newJobExecutor(1); + ThreadPoolExecutor secondJobExecutor = newJobExecutor(1); + WorkerJobPollingScheduler firstPoller = newPoller(firstWorkerId, firstJobExecutor, pipeline, 1); + WorkerJobPollingScheduler secondPoller = newPoller(secondWorkerId, secondJobExecutor, pipeline, 1); + + try { + // 1. 두 Poller가 같은 시점에 실제 Claim Transaction을 시작하도록 Barrier에서 맞춘다. + runPollersConcurrently(List.of(firstPoller, secondPoller)); + + // 2. 유일한 Claim만 Pipeline에 제출됐는지와 DB 소유권을 함께 확인한다. + assertThat(executed.await(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isTrue(); + then(pipeline).should(times(1)).execute( + any(ClaimedEmbeddingJobResponse.class), + any(WorkerExecutionSlot.class) + ); + assertThat(queryInteger( + "SELECT COUNT(*) FROM embedding_jobs WHERE id = ? AND status = 'PROCESSING'", + jobId + )).isOne(); + assertThat(queryLong( + "SELECT locked_by_worker_id FROM embedding_jobs WHERE id = ?", + jobId + )).isIn(firstWorkerId, secondWorkerId); + assertThat(eventCount(jobId, "LOCKED")).isOne(); + } finally { + shutdownExecutor(firstJobExecutor); + shutdownExecutor(secondJobExecutor); + } + } + + @Test + @DisplayName("한 Worker는 실행 슬롯 수보다 많은 Job을 PROCESSING으로 Claim하지 않는다") + void poll_limitsProcessingJobsToLocalExecutionCapacity() throws Exception { + TestContext context = insertContext("slot-capacity"); + insertPendingJobs(context, 5); + Long workerId = insertActiveWorker("capacity-poller"); + WorkerExecutionSlotPool slotPool = new WorkerExecutionSlotPool(MAX_CONCURRENCY); + WorkerIndexingPipeline pipeline = mock(WorkerIndexingPipeline.class); + CountDownLatch started = new CountDownLatch(MAX_CONCURRENCY); + CountDownLatch release = new CountDownLatch(1); + doAnswer(invocation -> { + WorkerExecutionSlot slot = invocation.getArgument(1); + try { + started.countDown(); + if (!release.await(TIMEOUT_SECONDS, TimeUnit.SECONDS)) { + throw new IllegalStateException("Worker Pipeline 해제 대기가 제한 시간을 초과했습니다."); + } + } catch (InterruptedException exception) { + Thread.currentThread().interrupt(); + throw new IllegalStateException("Worker Pipeline 대기 중 Thread가 중단되었습니다.", exception); + } finally { + slot.close(); + } + return null; + }).when(pipeline).execute(any(ClaimedEmbeddingJobResponse.class), any(WorkerExecutionSlot.class)); + + ThreadPoolExecutor jobExecutor = newJobExecutor(MAX_CONCURRENCY); + WorkerJobPollingScheduler poller = newPoller(workerId, jobExecutor, pipeline, slotPool); + + try { + // 1. 첫 Polling이 두 실행 슬롯을 모두 점유한 상태로 Pipeline을 대기시킨다. + poller.poll(); + assertThat(started.await(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isTrue(); + assertThat(slotPool.getActiveSlots()).isEqualTo(MAX_CONCURRENCY); + + // 2. 슬롯이 없는 동안 다시 Polling해도 추가 Claim이 발생하지 않는지 DB에서 검증한다. + poller.poll(); + assertThat(countJobs("PROCESSING", workerId)).isEqualTo(MAX_CONCURRENCY); + assertThat(countJobs("PENDING", null)).isEqualTo(3); + assertThat(countEvents("LOCKED")).isEqualTo(MAX_CONCURRENCY); + + // 3. 실행이 끝나면 모든 슬롯이 반환돼 다음 Polling이 가능해진다. + release.countDown(); + awaitAvailableSlots(slotPool, MAX_CONCURRENCY); + } finally { + release.countDown(); + shutdownExecutor(jobExecutor); + } + } + + @Test + @DisplayName("활성 실행의 Lease가 갱신되면 원래 만료 시각의 Recovery가 Job을 회수하지 않는다") + void leaseRenewal_preventsRecoveryAtOriginalExpiry() throws Exception { + TestContext context = insertContext("renewal-fencing"); + Long jobId = insertPendingJobs(context, 1).get(0); + Long workerId = insertActiveWorker("renewal-worker"); + ClaimedEmbeddingJobResponse claimedJob = claimService.claim(workerId).orElseThrow(); + LocalDateTime originalExpiry = claimedJob.lockExpiresAt(); + ScheduledThreadPoolExecutor leaseScheduler = new ScheduledThreadPoolExecutor(1); + WorkerLeaseRenewalManager renewalManager = new WorkerLeaseRenewalManager( + leaseScheduler, + leaseService, + workerProperties + ); + WorkerLeaseRenewalHandle handle = renewalManager.start(claimedJob); + + try { + // 1. 실제 Lease Service가 원래 만료 시각보다 뒤로 DB Lease를 연장할 때까지 기다린다. + LocalDateTime renewedExpiry = awaitRenewedExpiry(jobId, originalExpiry); + assertThat(renewedExpiry).isAfter(originalExpiry); + + // 2. 원래 Lease는 만료된 논리 시각이어도 갱신된 DB Lease가 Recovery를 차단해야 한다. + RecoveryResult result = recoveryService.recover(jobId, originalExpiry.plusNanos(1_000)); + assertThat(result.recovered()).isFalse(); + assertThat(queryString("SELECT status FROM embedding_jobs WHERE id = ?", jobId)) + .isEqualTo("PROCESSING"); + assertThat(eventCount(jobId, "LEASE_EXPIRED")).isZero(); + } finally { + handle.close(); + shutdownLeaseManager(renewalManager, leaseScheduler); + } + } + + @Test + @DisplayName("Lease 갱신을 중단한 뒤 만료 Job을 동시에 복구해도 한 번만 재예약한다") + void stoppedLeaseRenewal_allowsExactlyOneConcurrentRecovery() throws Exception { + TestContext context = insertContext("renewal-stopped"); + Long jobId = insertPendingJobs(context, 1).get(0); + Long workerId = insertActiveWorker("stopped-renewal-worker"); + ClaimedEmbeddingJobResponse claimedJob = claimService.claim(workerId).orElseThrow(); + ScheduledThreadPoolExecutor leaseScheduler = new ScheduledThreadPoolExecutor(1); + WorkerLeaseRenewalManager renewalManager = new WorkerLeaseRenewalManager( + leaseScheduler, + leaseService, + workerProperties + ); + WorkerLeaseRenewalHandle handle = renewalManager.start(claimedJob); + + // 1. 한 번 이상 실제 갱신한 뒤 Handle과 Scheduler를 닫아 이후 Lease 연장을 중단한다. + awaitRenewedExpiry(jobId, claimedJob.lockExpiresAt()); + handle.close(); + shutdownLeaseManager(renewalManager, leaseScheduler); + LocalDateTime finalExpiry = queryDateTime( + "SELECT lock_expires_at FROM embedding_jobs WHERE id = ?", + jobId + ); + + // 2. 같은 만료 시각에 두 Recovery Transaction을 경쟁시켜 유일한 상태 전이만 허용한다. + List results = recoverConcurrently(jobId, finalExpiry); + + // 3. 한 번의 Retry와 한 세트의 Event만 Commit됐는지 최종 DB 상태로 검증한다. + assertThat(results).filteredOn(RecoveryResult::recovered).hasSize(1); + assertThat(results).filteredOn(result -> !result.recovered()).hasSize(1); + assertThat(queryString("SELECT status FROM embedding_jobs WHERE id = ?", jobId)) + .isEqualTo("PENDING"); + assertThat(queryInteger("SELECT retry_count FROM embedding_jobs WHERE id = ?", jobId)).isOne(); + assertThat(queryDateTime("SELECT next_retry_at FROM embedding_jobs WHERE id = ?", jobId)) + .isAfter(finalExpiry); + assertThat(eventCount(jobId, "LEASE_EXPIRED")).isOne(); + assertThat(eventCount(jobId, "RETRY")).isOne(); + } + + private WorkerJobPollingScheduler newPoller( + Long workerId, + ThreadPoolExecutor jobExecutor, + WorkerIndexingPipeline pipeline, + int capacity + ) { + return newPoller(workerId, jobExecutor, pipeline, new WorkerExecutionSlotPool(capacity)); + } + + private WorkerJobPollingScheduler newPoller( + Long workerId, + ThreadPoolExecutor jobExecutor, + WorkerIndexingPipeline pipeline, + WorkerExecutionSlotPool slotPool + ) { + WorkerLifecycleManager lifecycleManager = mock(WorkerLifecycleManager.class); + given(lifecycleManager.getWorkerId()).willReturn(Optional.of(workerId)); + return new WorkerJobPollingScheduler( + lifecycleManager, + claimService, + slotPool, + jobExecutor, + pipeline + ); + } + + private ThreadPoolExecutor newJobExecutor(int capacity) { + return new ThreadPoolExecutor( + capacity, + capacity, + 0L, + TimeUnit.MILLISECONDS, + new SynchronousQueue<>() + ); + } + + private void runPollersConcurrently(List pollers) throws Exception { + CyclicBarrier startBarrier = new CyclicBarrier(pollers.size()); + ExecutorService executor = Executors.newFixedThreadPool(pollers.size()); + try { + List> futures = new ArrayList<>(); + for (WorkerJobPollingScheduler poller : pollers) { + futures.add(executor.submit(() -> { + startBarrier.await(TIMEOUT_SECONDS, TimeUnit.SECONDS); + poller.poll(); + return null; + })); + } + for (Future future : futures) { + future.get(TIMEOUT_SECONDS, TimeUnit.SECONDS); + } + } finally { + executor.shutdownNow(); + assertThat(executor.awaitTermination(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isTrue(); + } + } + + private List recoverConcurrently( + Long jobId, + LocalDateTime recoveredAt + ) throws Exception { + CyclicBarrier startBarrier = new CyclicBarrier(2); + ExecutorService executor = Executors.newFixedThreadPool(2); + try { + List> futures = List.of( + executor.submit(() -> recoverAfterBarrier(jobId, recoveredAt, startBarrier)), + executor.submit(() -> recoverAfterBarrier(jobId, recoveredAt, startBarrier)) + ); + return List.of( + futures.get(0).get(TIMEOUT_SECONDS, TimeUnit.SECONDS), + futures.get(1).get(TIMEOUT_SECONDS, TimeUnit.SECONDS) + ); + } finally { + executor.shutdownNow(); + assertThat(executor.awaitTermination(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isTrue(); + } + } + + private RecoveryResult recoverAfterBarrier( + Long jobId, + LocalDateTime recoveredAt, + CyclicBarrier startBarrier + ) throws Exception { + startBarrier.await(TIMEOUT_SECONDS, TimeUnit.SECONDS); + return recoveryService.recover(jobId, recoveredAt); + } + + private LocalDateTime awaitRenewedExpiry( + Long jobId, + LocalDateTime originalExpiry + ) throws InterruptedException { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(TIMEOUT_SECONDS); + while (System.nanoTime() < deadline) { + LocalDateTime currentExpiry = queryDateTime( + "SELECT lock_expires_at FROM embedding_jobs WHERE id = ?", + jobId + ); + if (currentExpiry.isAfter(originalExpiry)) { + return currentExpiry; + } + Thread.sleep(20L); + } + throw new IllegalStateException("Lease 갱신이 제한 시간 안에 DB에 반영되지 않았습니다."); + } + + private void awaitAvailableSlots( + WorkerExecutionSlotPool slotPool, + int expectedSlots + ) throws InterruptedException { + long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(TIMEOUT_SECONDS); + while (System.nanoTime() < deadline) { + if (slotPool.getAvailableSlots() == expectedSlots) { + return; + } + Thread.sleep(20L); + } + throw new IllegalStateException("Worker 실행 슬롯이 제한 시간 안에 반환되지 않았습니다."); + } + + private void shutdownExecutor(ThreadPoolExecutor executor) throws InterruptedException { + executor.shutdownNow(); + assertThat(executor.awaitTermination(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isTrue(); + } + + private void shutdownLeaseManager( + WorkerLeaseRenewalManager renewalManager, + ScheduledThreadPoolExecutor leaseScheduler + ) throws InterruptedException { + renewalManager.stopAll(); + leaseScheduler.shutdownNow(); + assertThat(leaseScheduler.awaitTermination(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isTrue(); + } + + private TestContext insertContext(String scenario) { + String suffix = UUID.randomUUID().toString(); + Long userId = jdbcTemplate.queryForObject(""" + INSERT INTO users (email, password_hash, name, status, created_at, updated_at) + VALUES (?, 'password-hash', 'Worker Orchestration User', 'ACTIVE', + CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) + RETURNING id + """, Long.class, "worker-orchestration-" + suffix + "@example.com"); + Long documentId = jdbcTemplate.queryForObject(""" + INSERT INTO documents ( + owner_user_id, title, document_type, source_type, status, visibility, + created_at, updated_at + ) + VALUES (?, ?, 'TXT', 'UPLOAD', 'UPLOADED', 'PRIVATE', + CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) + RETURNING id + """, Long.class, userId, "Worker Orchestration " + scenario + " " + suffix); + Long versionId = jdbcTemplate.queryForObject(""" + INSERT INTO document_versions ( + document_id, version_no, title_snapshot, content_type, status, + created_by, created_at, updated_at + ) + VALUES (?, 1, ?, 'text/plain', 'PARSING', ?, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) + RETURNING id + """, Long.class, documentId, "Worker Orchestration " + scenario, userId); + jdbcTemplate.update( + "UPDATE documents SET current_version_id = ? WHERE id = ?", + versionId, + documentId + ); + Long embeddingModelId = jdbcTemplate.queryForObject(""" + SELECT id + FROM embedding_models + WHERE is_active = TRUE AND is_searchable = TRUE + """, Long.class); + return new TestContext(versionId, embeddingModelId); + } + + private List insertPendingJobs(TestContext context, int count) { + List jobIds = new ArrayList<>(); + for (int index = 0; index < count; index++) { + jobIds.add(jdbcTemplate.queryForObject(""" + INSERT INTO embedding_jobs ( + document_version_id, embedding_model_id, status, priority, + retry_count, max_retry_count, created_at, updated_at + ) + VALUES (?, ?, 'PENDING', 0, 0, 3, + CURRENT_TIMESTAMP + ? * INTERVAL '1 microsecond', CURRENT_TIMESTAMP) + RETURNING id + """, Long.class, context.versionId(), context.embeddingModelId(), index)); + } + return jobIds; + } + + private Long insertActiveWorker(String workerName) { + return jdbcTemplate.queryForObject(""" + INSERT INTO worker_nodes ( + worker_name, instance_id, host_name, ip_address, status, + last_heartbeat_at, started_at, created_at, updated_at + ) + VALUES (?, ?, 'localhost', '127.0.0.1', 'ACTIVE', CURRENT_TIMESTAMP, + CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) + RETURNING id + """, Long.class, workerName, UUID.randomUUID().toString()); + } + + private int countJobs(String status, Long workerId) { + if (workerId == null) { + return queryInteger("SELECT COUNT(*) FROM embedding_jobs WHERE status = ?", status); + } + return queryInteger( + "SELECT COUNT(*) FROM embedding_jobs WHERE status = ? AND locked_by_worker_id = ?", + status, + workerId + ); + } + + private int countEvents(String eventType) { + return queryInteger("SELECT COUNT(*) FROM indexing_events WHERE event_type = ?", eventType); + } + + private int eventCount(Long jobId, String eventType) { + return queryInteger( + "SELECT COUNT(*) FROM indexing_events WHERE embedding_job_id = ? AND event_type = ?", + jobId, + eventType + ); + } + + private int queryInteger(String sql, Object... arguments) { + return jdbcTemplate.queryForObject(sql, Integer.class, arguments); + } + + private Long queryLong(String sql, Object... arguments) { + return jdbcTemplate.queryForObject(sql, Long.class, arguments); + } + + private String queryString(String sql, Object... arguments) { + return jdbcTemplate.queryForObject(sql, String.class, arguments); + } + + private LocalDateTime queryDateTime(String sql, Object... arguments) { + return jdbcTemplate.queryForObject(sql, LocalDateTime.class, arguments); + } + + /** + * 한 테스트 문서와 Version, 활성 Embedding Model의 식별자를 전달하는 DB 준비 결과다. + */ + private record TestContext(Long versionId, Long embeddingModelId) { + } +} diff --git a/src/test/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerExecutionLifecycleManagerTest.java b/src/test/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerExecutionLifecycleManagerTest.java new file mode 100644 index 0000000..8f46615 --- /dev/null +++ b/src/test/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerExecutionLifecycleManagerTest.java @@ -0,0 +1,90 @@ +package com.opensource.docgrid.domain.worker.lifecycle; + +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.then; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.times; + +import java.time.Duration; +import java.util.List; +import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import com.opensource.docgrid.domain.worker.config.IndexingWorkerProperties; +import com.opensource.docgrid.domain.worker.service.WorkerLeaseRenewalManager; + +/** + * Worker 종료 시 Polling 차단, 실행 유예, 강제 중단과 Lease Scheduler 정리 순서를 검증한다. + */ +@ExtendWith(MockitoExtension.class) +@DisplayName("WorkerExecutionLifecycleManager 테스트") +class WorkerExecutionLifecycleManagerTest { + + @Mock private WorkerJobPollingScheduler pollingScheduler; + @Mock private ThreadPoolExecutor jobExecutor; + @Mock private ScheduledThreadPoolExecutor leaseScheduler; + @Mock private WorkerLeaseRenewalManager leaseRenewalManager; + + private WorkerExecutionLifecycleManager lifecycleManager; + + @BeforeEach + void setUp() { + IndexingWorkerProperties workerProperties = new IndexingWorkerProperties(); + workerProperties.setShutdownGracePeriod(Duration.ofSeconds(5)); + lifecycleManager = new WorkerExecutionLifecycleManager( + pollingScheduler, + jobExecutor, + leaseScheduler, + leaseRenewalManager, + workerProperties + ); + } + + @Test + @DisplayName("실행이 유예 시간 안에 끝나면 interrupt 없이 Lease 자원을 정리한다") + void shutdown_completesGracefully_whenExecutorTerminates() throws InterruptedException { + given(jobExecutor.awaitTermination(anyLong(), eq(TimeUnit.NANOSECONDS))) + .willReturn(true); + + lifecycleManager.shutdown(); + lifecycleManager.shutdown(); + + then(pollingScheduler).should().stopPolling(); + then(jobExecutor).should().shutdown(); + then(jobExecutor).should().awaitTermination( + Duration.ofSeconds(5).toNanos(), + TimeUnit.NANOSECONDS + ); + then(jobExecutor).should(never()).shutdownNow(); + then(leaseRenewalManager).should().stopAll(); + then(leaseScheduler).should().shutdown(); + then(pollingScheduler).should(times(1)).stopPolling(); + } + + @Test + @DisplayName("유예 시간이 끝나면 남은 실행을 interrupt한 뒤 Lease 자원을 정리한다") + void shutdown_interruptsJobs_whenGracePeriodExpires() throws InterruptedException { + given(jobExecutor.awaitTermination(anyLong(), eq(TimeUnit.NANOSECONDS))) + .willReturn(false); + given(jobExecutor.shutdownNow()).willReturn(List.of(() -> { + })); + + lifecycleManager.shutdown(); + + then(pollingScheduler).should().stopPolling(); + then(jobExecutor).should().shutdown(); + then(jobExecutor).should().shutdownNow(); + then(leaseRenewalManager).should().stopAll(); + then(leaseScheduler).should().shutdown(); + } +} diff --git a/src/test/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerJobPollingSchedulerTest.java b/src/test/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerJobPollingSchedulerTest.java new file mode 100644 index 0000000..a48389e --- /dev/null +++ b/src/test/java/com/opensource/docgrid/domain/worker/lifecycle/WorkerJobPollingSchedulerTest.java @@ -0,0 +1,157 @@ +package com.opensource.docgrid.domain.worker.lifecycle; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.then; +import static org.mockito.Mockito.times; + +import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.List; +import java.util.Optional; +import java.util.concurrent.RejectedExecutionException; +import java.util.concurrent.ThreadPoolExecutor; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import com.opensource.docgrid.domain.embedding.dto.response.ClaimedEmbeddingJobResponse; +import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus; +import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobClaimService; +import com.opensource.docgrid.domain.worker.execution.WorkerExecutionSlotPool; +import com.opensource.docgrid.domain.worker.execution.WorkerExecutionSlotPool.WorkerExecutionSlot; +import com.opensource.docgrid.domain.worker.service.WorkerIndexingPipeline; + +/** + * Worker Poller의 등록 조건, 실행 슬롯 제한, Claim 중단과 제출 거부 시 자원 반환을 검증한다. + */ +@ExtendWith(MockitoExtension.class) +@DisplayName("WorkerJobPollingScheduler 테스트") +class WorkerJobPollingSchedulerTest { + + @Mock private WorkerLifecycleManager workerLifecycleManager; + @Mock private EmbeddingJobClaimService claimService; + @Mock private ThreadPoolExecutor jobExecutor; + @Mock private WorkerIndexingPipeline indexingPipeline; + + private WorkerExecutionSlotPool slotPool; + private WorkerJobPollingScheduler scheduler; + + @BeforeEach + void setUp() { + slotPool = new WorkerExecutionSlotPool(2); + scheduler = new WorkerJobPollingScheduler( + workerLifecycleManager, + claimService, + slotPool, + jobExecutor, + indexingPipeline + ); + } + + @Test + @DisplayName("Worker 등록 전에는 Job을 Claim하지 않는다") + void poll_doesNotClaim_beforeWorkerRegistration() { + given(workerLifecycleManager.getWorkerId()).willReturn(Optional.empty()); + + scheduler.poll(); + + then(claimService).shouldHaveNoInteractions(); + then(jobExecutor).shouldHaveNoInteractions(); + assertThat(slotPool.getAvailableSlots()).isEqualTo(2); + } + + @Test + @DisplayName("빈 Claim은 획득한 Slot을 반환하고 현재 Polling을 끝낸다") + void poll_releasesSlot_whenNoJobIsClaimed() { + given(workerLifecycleManager.getWorkerId()).willReturn(Optional.of(1L)); + given(claimService.claim(1L)).willReturn(Optional.empty()); + + scheduler.poll(); + + then(claimService).should().claim(1L); + then(jobExecutor).shouldHaveNoInteractions(); + assertThat(slotPool.getAvailableSlots()).isEqualTo(2); + } + + @Test + @DisplayName("한 Polling 주기에는 실행 가능한 Slot 수만큼만 Claim해 제출한다") + void poll_claimsOnlyUpToAvailableCapacity() { + List submittedTasks = new ArrayList<>(); + given(workerLifecycleManager.getWorkerId()).willReturn(Optional.of(1L)); + given(claimService.claim(1L)) + .willReturn(Optional.of(claimedJob(10L))) + .willReturn(Optional.of(claimedJob(11L))); + org.mockito.Mockito.doAnswer(invocation -> { + submittedTasks.add(invocation.getArgument(0)); + return null; + }).when(jobExecutor).execute(any(Runnable.class)); + org.mockito.Mockito.doAnswer(invocation -> { + WorkerExecutionSlot slot = invocation.getArgument(1); + slot.close(); + return null; + }).when(indexingPipeline).execute( + any(ClaimedEmbeddingJobResponse.class), + any(WorkerExecutionSlot.class) + ); + + scheduler.poll(); + + then(claimService).should(times(2)).claim(1L); + assertThat(submittedTasks).hasSize(2); + assertThat(slotPool.getActiveSlots()).isEqualTo(2); + + submittedTasks.forEach(Runnable::run); + + then(indexingPipeline).should(times(2)).execute( + any(ClaimedEmbeddingJobResponse.class), + any(WorkerExecutionSlot.class) + ); + assertThat(slotPool.getAvailableSlots()).isEqualTo(2); + } + + @Test + @DisplayName("Executor가 제출을 거부하면 Slot을 반환하고 추가 Claim을 중단한다") + void poll_releasesSlot_whenSubmissionIsRejected() { + given(workerLifecycleManager.getWorkerId()).willReturn(Optional.of(1L)); + given(claimService.claim(1L)).willReturn(Optional.of(claimedJob(10L))); + org.mockito.Mockito.doThrow(new RejectedExecutionException("closing")) + .when(jobExecutor).execute(any(Runnable.class)); + + scheduler.poll(); + + then(claimService).should().claim(1L); + then(indexingPipeline).shouldHaveNoInteractions(); + assertThat(slotPool.getAvailableSlots()).isEqualTo(2); + } + + @Test + @DisplayName("종료 시작 후에는 신규 Polling과 Claim을 허용하지 않는다") + void stopPolling_blocksFollowingPolls() { + scheduler.stopPolling(); + + scheduler.poll(); + + assertThat(scheduler.isPollingEnabled()).isFalse(); + then(workerLifecycleManager).shouldHaveNoInteractions(); + then(claimService).shouldHaveNoInteractions(); + } + + private ClaimedEmbeddingJobResponse claimedJob(Long jobId) { + return new ClaimedEmbeddingJobResponse( + jobId, + EmbeddingJobStatus.PROCESSING, + 1L, + 5L, + 7L, + "34c19d16-6ae1-4f6a-a35d-0123456789ab", + LocalDateTime.of(2026, 8, 3, 18, 0), + LocalDateTime.of(2026, 8, 3, 18, 5) + ); + } +} diff --git a/src/test/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailureClassifierTest.java b/src/test/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailureClassifierTest.java new file mode 100644 index 0000000..9edaa38 --- /dev/null +++ b/src/test/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailureClassifierTest.java @@ -0,0 +1,70 @@ +package com.opensource.docgrid.domain.worker.service; + +import static org.assertj.core.api.Assertions.assertThat; + +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import com.opensource.docgrid.domain.embedding.enums.IndexingFailureType; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +/** + * Worker 예외를 제한 실패 유형과 고정 메시지로 분류하고 소유권 상실 보고를 차단하는지 검증한다. + */ +@DisplayName("WorkerIndexingFailureClassifier 테스트") +class WorkerIndexingFailureClassifierTest { + + private final WorkerIndexingFailureClassifier classifier = new WorkerIndexingFailureClassifier(); + + @Test + @DisplayName("저장소와 Provider 장애는 Retry 가능한 외부 실패로 분류한다") + void classify_mapsRetryableExternalFailures() { + WorkerIndexingFailure storage = classifier.classify( + new DocGridException(ErrorCode.FILE_STORAGE_FAILED, "sensitive object key") + ); + WorkerIndexingFailure provider = classifier.classify( + new DocGridException(ErrorCode.EMBEDDING_SERVER_UNAVAILABLE, "sensitive endpoint") + ); + + assertThat(storage.failureType()).isEqualTo(IndexingFailureType.STORAGE_UNAVAILABLE); + assertThat(storage.safeMessage()).doesNotContain("sensitive"); + assertThat(provider.failureType()) + .isEqualTo(IndexingFailureType.EMBEDDING_PROVIDER_UNAVAILABLE); + assertThat(provider.safeMessage()).doesNotContain("sensitive"); + } + + @Test + @DisplayName("문서·Vector·상태 오류를 영구 실패 유형으로 분류한다") + void classify_mapsNonRetryableFailures() { + assertThat(classifier.classify(new DocGridException(ErrorCode.DOCUMENT_CONTENT_EMPTY)) + .failureType()).isEqualTo(IndexingFailureType.DOCUMENT_CONTENT_INVALID); + assertThat(classifier.classify(new DocGridException(ErrorCode.EMBEDDING_VECTOR_INVALID)) + .failureType()).isEqualTo(IndexingFailureType.EMBEDDING_RESULT_INVALID); + assertThat(classifier.classify(new DocGridException(ErrorCode.DOCUMENT_CHUNKS_INCONSISTENT)) + .failureType()).isEqualTo(IndexingFailureType.INDEXING_STATE_INCONSISTENT); + } + + @Test + @DisplayName("소유권과 Lease 오류는 과거 Claim 실패 보고에서 제외한다") + void classify_skipsOwnershipLostFailures() { + WorkerIndexingFailure failure = classifier.classify( + new DocGridException(ErrorCode.EMBEDDING_JOB_LEASE_EXPIRED) + ); + + assertThat(failure.reportable()).isFalse(); + assertThat(failure.failureType()).isNull(); + assertThat(failure.diagnosticCode()).isEqualTo("EMBEDDING-JOB-004"); + } + + @Test + @DisplayName("알 수 없는 Runtime 오류는 Worker 내부 Retry 유형으로 제한한다") + void classify_mapsUnknownRuntimeFailure() { + WorkerIndexingFailure failure = classifier.classify( + new IllegalStateException("sensitive runtime detail") + ); + + assertThat(failure.failureType()).isEqualTo(IndexingFailureType.WORKER_INTERNAL_ERROR); + assertThat(failure.safeMessage()).doesNotContain("sensitive"); + } +} diff --git a/src/test/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailureReporterTest.java b/src/test/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailureReporterTest.java new file mode 100644 index 0000000..7b678ce --- /dev/null +++ b/src/test/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingFailureReporterTest.java @@ -0,0 +1,107 @@ +package com.opensource.docgrid.domain.worker.service; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.then; +import static org.mockito.Mockito.never; + +import java.time.LocalDateTime; + +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import com.opensource.docgrid.domain.embedding.dto.request.FailDocumentIndexingRequest; +import com.opensource.docgrid.domain.embedding.dto.response.ClaimedEmbeddingJobResponse; +import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus; +import com.opensource.docgrid.domain.embedding.enums.IndexingFailureType; +import com.opensource.docgrid.domain.embedding.service.command.DocumentIndexingFailureService; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +/** + * Worker 실패 Reporter가 고정된 분류 결과만 전달하고 소유권 상실은 보고하지 않는지 검증한다. + */ +@ExtendWith(MockitoExtension.class) +@DisplayName("WorkerIndexingFailureReporter 테스트") +class WorkerIndexingFailureReporterTest { + + @Mock private DocumentIndexingFailureService failureService; + + @Test + @DisplayName("분류된 유형과 안전한 메시지로 현재 Attempt 실패를 보고한다") + void report_sendsClassifiedFailure() { + WorkerIndexingFailureReporter reporter = new WorkerIndexingFailureReporter( + new WorkerIndexingFailureClassifier(), + failureService + ); + ClaimedEmbeddingJobResponse claimedJob = claimedJob(); + ArgumentCaptor requestCaptor = + ArgumentCaptor.forClass(FailDocumentIndexingRequest.class); + + reporter.report( + claimedJob, + 100L, + new DocGridException(ErrorCode.FILE_STORAGE_FAILED, "sensitive object key") + ); + + then(failureService).should().fail( + org.mockito.ArgumentMatchers.eq(10L), + org.mockito.ArgumentMatchers.eq(100L), + requestCaptor.capture() + ); + assertThat(requestCaptor.getValue().failureType()) + .isEqualTo(IndexingFailureType.STORAGE_UNAVAILABLE); + assertThat(requestCaptor.getValue().errorMessage()).doesNotContain("sensitive"); + assertThat(requestCaptor.getValue().claimToken()).isEqualTo(claimedJob.claimToken()); + } + + @Test + @DisplayName("소유권을 잃었으면 실패 Service를 호출하지 않는다") + void report_skipsOwnershipLostFailure() { + WorkerIndexingFailureReporter reporter = new WorkerIndexingFailureReporter( + new WorkerIndexingFailureClassifier(), + failureService + ); + + reporter.report( + claimedJob(), + 100L, + new DocGridException(ErrorCode.EMBEDDING_JOB_OWNERSHIP_INVALID) + ); + + then(failureService).should(never()).fail(any(), any(), any()); + } + + @Test + @DisplayName("실패 보고 자체가 실패해도 예외를 외부로 전파하지 않는다") + void report_isolatesReportingFailure() { + WorkerIndexingFailureReporter reporter = new WorkerIndexingFailureReporter( + new WorkerIndexingFailureClassifier(), + failureService + ); + given(failureService.fail(any(), any(), any())) + .willThrow(new IllegalStateException("report failed")); + + reporter.report(claimedJob(), 100L, new IllegalStateException("pipeline failed")); + + then(failureService).should().fail(any(), any(), any()); + } + + private ClaimedEmbeddingJobResponse claimedJob() { + return new ClaimedEmbeddingJobResponse( + 10L, + EmbeddingJobStatus.PROCESSING, + 1L, + 5L, + 7L, + "34c19d16-6ae1-4f6a-a35d-0123456789ab", + LocalDateTime.of(2026, 8, 3, 18, 0), + LocalDateTime.of(2026, 8, 3, 18, 5) + ); + } +} diff --git a/src/test/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingPipelineTest.java b/src/test/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingPipelineTest.java new file mode 100644 index 0000000..9d073f7 --- /dev/null +++ b/src/test/java/com/opensource/docgrid/domain/worker/service/WorkerIndexingPipelineTest.java @@ -0,0 +1,215 @@ +package com.opensource.docgrid.domain.worker.service; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.then; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.never; + +import java.time.LocalDateTime; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.InOrder; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import com.opensource.docgrid.domain.document.enums.DocumentVersionStatus; +import com.opensource.docgrid.domain.document.service.DocumentParsingService; +import com.opensource.docgrid.domain.document.service.query.DocumentIndexingStageQueryService; +import com.opensource.docgrid.domain.embedding.dto.request.CompleteDocumentIndexingRequest; +import com.opensource.docgrid.domain.embedding.dto.request.CreateDocumentChunksRequest; +import com.opensource.docgrid.domain.embedding.dto.request.CreateDocumentEmbeddingsRequest; +import com.opensource.docgrid.domain.embedding.dto.request.StartEmbeddingJobAttemptRequest; +import com.opensource.docgrid.domain.embedding.dto.response.ClaimedEmbeddingJobResponse; +import com.opensource.docgrid.domain.embedding.dto.response.StartedEmbeddingJobAttemptResponse; +import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus; +import com.opensource.docgrid.domain.embedding.service.DocumentEmbeddingService; +import com.opensource.docgrid.domain.embedding.service.command.DocumentIndexingCompletionService; +import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobAttemptService; +import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobAttemptService.StartResult; +import com.opensource.docgrid.domain.worker.enums.AttemptStatus; +import com.opensource.docgrid.domain.worker.execution.WorkerExecutionSlotPool; +import com.opensource.docgrid.domain.worker.execution.WorkerExecutionSlotPool.WorkerExecutionSlot; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +/** + * Worker Pipeline의 상태별 단계 선택, Service 순서와 실패·자원 정리 경계를 검증한다. + */ +@ExtendWith(MockitoExtension.class) +@DisplayName("WorkerIndexingPipeline 테스트") +class WorkerIndexingPipelineTest { + + private static final Long JOB_ID = 10L; + private static final Long ATTEMPT_ID = 100L; + private static final Long WORKER_ID = 1L; + private static final Long VERSION_ID = 5L; + private static final String CLAIM_TOKEN = "34c19d16-6ae1-4f6a-a35d-0123456789ab"; + + @Mock private EmbeddingJobAttemptService attemptService; + @Mock private DocumentIndexingStageQueryService stageQueryService; + @Mock private DocumentParsingService parsingService; + @Mock private DocumentEmbeddingService embeddingService; + @Mock private DocumentIndexingCompletionService completionService; + @Mock private WorkerIndexingFailureReporter failureReporter; + @Mock private WorkerLeaseRenewalManager leaseRenewalManager; + + private WorkerIndexingPipeline pipeline; + private ClaimedEmbeddingJobResponse claimedJob; + + @BeforeEach + void setUp() { + pipeline = new WorkerIndexingPipeline( + attemptService, + stageQueryService, + parsingService, + embeddingService, + completionService, + failureReporter, + leaseRenewalManager + ); + claimedJob = claimedJob(); + } + + @Test + @DisplayName("UPLOADED 상태는 Attempt, Chunk, Embedding, 완료 순서로 실행한다") + void execute_runsAllStages_fromUploaded() { + given(attemptService.start(any(), any(StartEmbeddingJobAttemptRequest.class))) + .willReturn(startResult()); + WorkerLeaseRenewalHandle leaseHandle = leaseHandle(); + given(leaseRenewalManager.start(claimedJob)).willReturn(leaseHandle); + given(stageQueryService.getStatus(VERSION_ID)).willReturn(DocumentVersionStatus.UPLOADED); + WorkerExecutionSlotPool slotPool = new WorkerExecutionSlotPool(1); + + pipeline.execute(claimedJob, slotPool.tryAcquire().orElseThrow()); + + InOrder order = inOrder( + attemptService, + stageQueryService, + parsingService, + embeddingService, + completionService + ); + order.verify(attemptService).start(any(), any(StartEmbeddingJobAttemptRequest.class)); + order.verify(stageQueryService).getStatus(VERSION_ID); + order.verify(parsingService).createChunks( + org.mockito.ArgumentMatchers.eq(JOB_ID), + org.mockito.ArgumentMatchers.eq(ATTEMPT_ID), + any(CreateDocumentChunksRequest.class) + ); + order.verify(embeddingService).createEmbeddings( + org.mockito.ArgumentMatchers.eq(JOB_ID), + org.mockito.ArgumentMatchers.eq(ATTEMPT_ID), + any(CreateDocumentEmbeddingsRequest.class) + ); + order.verify(completionService).complete( + org.mockito.ArgumentMatchers.eq(JOB_ID), + org.mockito.ArgumentMatchers.eq(ATTEMPT_ID), + any(CompleteDocumentIndexingRequest.class) + ); + assertThat(slotPool.getAvailableSlots()).isOne(); + assertThat(leaseHandle.isClosed()).isTrue(); + then(failureReporter).shouldHaveNoInteractions(); + } + + @Test + @DisplayName("CHUNKED 상태는 Chunk를 건너뛰고 Embedding부터 실행한다") + void execute_skipsChunks_fromChunked() { + given(attemptService.start(any(), any(StartEmbeddingJobAttemptRequest.class))) + .willReturn(startResult()); + given(leaseRenewalManager.start(claimedJob)).willReturn(leaseHandle()); + given(stageQueryService.getStatus(VERSION_ID)).willReturn(DocumentVersionStatus.CHUNKED); + + pipeline.execute(claimedJob, slot()); + + then(parsingService).shouldHaveNoInteractions(); + then(embeddingService).should().createEmbeddings( + org.mockito.ArgumentMatchers.eq(JOB_ID), + org.mockito.ArgumentMatchers.eq(ATTEMPT_ID), + any(CreateDocumentEmbeddingsRequest.class) + ); + then(completionService).should().complete( + org.mockito.ArgumentMatchers.eq(JOB_ID), + org.mockito.ArgumentMatchers.eq(ATTEMPT_ID), + any(CompleteDocumentIndexingRequest.class) + ); + } + + @Test + @DisplayName("Attempt 시작 후 단계 오류는 실패 Reporter에 전달하고 Slot과 Lease를 닫는다") + void execute_reportsFailure_afterAttemptStart() { + given(attemptService.start(any(), any(StartEmbeddingJobAttemptRequest.class))) + .willReturn(startResult()); + WorkerLeaseRenewalHandle leaseHandle = leaseHandle(); + WorkerExecutionSlotPool slotPool = new WorkerExecutionSlotPool(1); + given(leaseRenewalManager.start(claimedJob)).willReturn(leaseHandle); + DocGridException failure = new DocGridException(ErrorCode.DOCUMENT_CHUNKS_INCONSISTENT); + given(stageQueryService.getStatus(VERSION_ID)).willThrow(failure); + + pipeline.execute(claimedJob, slotPool.tryAcquire().orElseThrow()); + + then(failureReporter).should().report(claimedJob, ATTEMPT_ID, failure); + then(embeddingService).shouldHaveNoInteractions(); + then(completionService).shouldHaveNoInteractions(); + assertThat(leaseHandle.isClosed()).isTrue(); + assertThat(slotPool.getAvailableSlots()).isOne(); + } + + @Test + @DisplayName("Attempt 시작 전 오류는 실패를 합성하지 않고 Slot만 반환한다") + void execute_propagatesFailure_beforeAttemptStart() { + WorkerExecutionSlotPool slotPool = new WorkerExecutionSlotPool(1); + given(attemptService.start(any(), any(StartEmbeddingJobAttemptRequest.class))) + .willThrow(new DocGridException(ErrorCode.EMBEDDING_JOB_LEASE_EXPIRED)); + + assertThatThrownBy(() -> pipeline.execute( + claimedJob, + slotPool.tryAcquire().orElseThrow() + )).isInstanceOf(DocGridException.class); + + then(leaseRenewalManager).should(never()).start(any()); + then(failureReporter).shouldHaveNoInteractions(); + assertThat(slotPool.getAvailableSlots()).isOne(); + } + + private WorkerExecutionSlot slot() { + return new WorkerExecutionSlotPool(1).tryAcquire().orElseThrow(); + } + + private WorkerLeaseRenewalHandle leaseHandle() { + return new WorkerLeaseRenewalHandle(JOB_ID, WORKER_ID, CLAIM_TOKEN, ignored -> { + }); + } + + private StartResult startResult() { + return new StartResult( + new StartedEmbeddingJobAttemptResponse( + ATTEMPT_ID, + JOB_ID, + 1, + WORKER_ID, + AttemptStatus.STARTED, + LocalDateTime.of(2026, 8, 3, 18, 0) + ), + true + ); + } + + private ClaimedEmbeddingJobResponse claimedJob() { + return new ClaimedEmbeddingJobResponse( + JOB_ID, + EmbeddingJobStatus.PROCESSING, + WORKER_ID, + VERSION_ID, + 7L, + CLAIM_TOKEN, + LocalDateTime.of(2026, 8, 3, 18, 0), + LocalDateTime.of(2026, 8, 3, 18, 5) + ); + } +} diff --git a/src/test/java/com/opensource/docgrid/domain/worker/service/WorkerLeaseRenewalManagerTest.java b/src/test/java/com/opensource/docgrid/domain/worker/service/WorkerLeaseRenewalManagerTest.java new file mode 100644 index 0000000..b7bb1ee --- /dev/null +++ b/src/test/java/com/opensource/docgrid/domain/worker/service/WorkerLeaseRenewalManagerTest.java @@ -0,0 +1,140 @@ +package com.opensource.docgrid.domain.worker.service; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatCode; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.then; + +import java.time.Duration; +import java.time.LocalDateTime; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.ScheduledThreadPoolExecutor; +import java.util.concurrent.TimeUnit; + +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import com.opensource.docgrid.domain.embedding.dto.request.RenewEmbeddingJobLeaseRequest; +import com.opensource.docgrid.domain.embedding.dto.response.ClaimedEmbeddingJobResponse; +import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus; +import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobLeaseService; +import com.opensource.docgrid.domain.worker.config.IndexingWorkerProperties; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +/** + * 실행별 Lease 갱신 예약, 소유권 상실, 일시 오류 유지와 전체 종료 동작을 검증한다. + */ +@ExtendWith(MockitoExtension.class) +@DisplayName("WorkerLeaseRenewalManager 테스트") +class WorkerLeaseRenewalManagerTest { + + @Mock private ScheduledThreadPoolExecutor leaseScheduler; + @Mock private ScheduledFuture scheduledFuture; + @Mock private EmbeddingJobLeaseService leaseService; + + private IndexingWorkerProperties workerProperties; + private WorkerLeaseRenewalManager manager; + private ArgumentCaptor taskCaptor; + + @BeforeEach + void setUp() { + workerProperties = new IndexingWorkerProperties(); + workerProperties.setLeaseRenewalInterval(Duration.ofSeconds(1)); + manager = new WorkerLeaseRenewalManager(leaseScheduler, leaseService, workerProperties); + taskCaptor = ArgumentCaptor.forClass(Runnable.class); + org.mockito.Mockito.doReturn(scheduledFuture).when(leaseScheduler).scheduleWithFixedDelay( + taskCaptor.capture(), + anyLong(), + anyLong(), + eq(TimeUnit.NANOSECONDS) + ); + } + + @Test + @DisplayName("설정 주기로 현재 Claim Lease를 갱신하고 close 시 예약을 제거한다") + void start_schedulesRenewalAndCloseRemovesHandle() { + ClaimedEmbeddingJobResponse claimedJob = claimedJob(); + WorkerLeaseRenewalHandle handle = manager.start(claimedJob); + ArgumentCaptor requestCaptor = + ArgumentCaptor.forClass(RenewEmbeddingJobLeaseRequest.class); + + taskCaptor.getValue().run(); + + then(leaseService).should().renew(eq(10L), requestCaptor.capture()); + assertThat(requestCaptor.getValue().workerId()).isEqualTo(1L); + assertThat(requestCaptor.getValue().claimToken()).isEqualTo(claimedJob.claimToken()); + assertThat(manager.getActiveHandleCount()).isOne(); + + handle.close(); + then(scheduledFuture).should().cancel(false); + assertThat(manager.getActiveHandleCount()).isZero(); + } + + @Test + @DisplayName("권위 있는 갱신 거부는 Handle을 lost로 전환한다") + void renew_marksOwnershipLost_onOwnershipFailure() { + given(leaseService.renew(any(), any())) + .willThrow(new DocGridException(ErrorCode.EMBEDDING_JOB_OWNERSHIP_INVALID)); + WorkerLeaseRenewalHandle handle = manager.start(claimedJob()); + + taskCaptor.getValue().run(); + + assertThatThrownBy(handle::ensureOwned) + .isInstanceOfSatisfying(DocGridException.class, + exception -> assertThat(exception.getErrorCode()) + .isEqualTo(ErrorCode.EMBEDDING_JOB_OWNERSHIP_INVALID)); + assertThat(manager.getActiveHandleCount()).isZero(); + then(scheduledFuture).should().cancel(false); + } + + @Test + @DisplayName("일시 Runtime 오류는 Handle을 유지해 다음 갱신 기회를 남긴다") + void renew_keepsHandle_onTransientRuntimeFailure() { + given(leaseService.renew(any(), any())) + .willThrow(new IllegalStateException("temporary database failure")); + WorkerLeaseRenewalHandle handle = manager.start(claimedJob()); + + taskCaptor.getValue().run(); + + assertThatCode(handle::ensureOwned).doesNotThrowAnyException(); + assertThat(manager.getActiveHandleCount()).isOne(); + handle.close(); + } + + @Test + @DisplayName("전체 종료는 활성 예약을 닫고 새 Handle 등록을 거부한다") + void stopAll_closesHandlesAndRejectsNewOnes() { + WorkerLeaseRenewalHandle handle = manager.start(claimedJob()); + + manager.stopAll(); + + assertThat(handle.isClosed()).isTrue(); + assertThat(manager.isAccepting()).isFalse(); + assertThat(manager.getActiveHandleCount()).isZero(); + assertThatThrownBy(() -> manager.start(claimedJob())) + .isInstanceOf(IllegalStateException.class); + } + + private ClaimedEmbeddingJobResponse claimedJob() { + return new ClaimedEmbeddingJobResponse( + 10L, + EmbeddingJobStatus.PROCESSING, + 1L, + 5L, + 7L, + "34c19d16-6ae1-4f6a-a35d-0123456789ab", + LocalDateTime.of(2026, 8, 3, 18, 0), + LocalDateTime.of(2026, 8, 3, 18, 5) + ); + } +}