diff --git a/docs/design/Gimini-3-#82-chunk-embedding-vector-storage.md b/docs/design/Gimini-3-#82-chunk-embedding-vector-storage.md
new file mode 100644
index 0000000..1680492
--- /dev/null
+++ b/docs/design/Gimini-3-#82-chunk-embedding-vector-storage.md
@@ -0,0 +1,406 @@
+# Chunk Embedding 생성 및 Vector 저장 설계
+
+## 1. 목표
+
+현재 Worker가 소유한 `PROCESSING` Embedding Job의 `CHUNKED` Document Version을 대상으로,
+Job에 고정된 Embedding Model을 사용해 모든 Chunk의 Vector를 생성하고 `embeddings`에 원자적으로
+저장한다.
+
+핵심 결과는 다음과 같다.
+
+- Job ID, Attempt ID, Worker ID, Claim Token을 하나의 실행 Context로 유지한다.
+- 현재 활성 모델을 다시 선택하지 않고 `EmbeddingJob.embeddingModel`을 사용한다.
+- `chunk_index` 오름차순으로 Chunk를 한 건씩 외부 Embedding Server에 전달한다.
+- 외부 HTTP 호출 동안 DB Transaction과 행 잠금을 유지하지 않는다.
+- Vector 전체가 검증된 경우에만 Version·Model 단위 Embedding Set을 한 Transaction으로 저장한다.
+- 처리 전후에 현재 Claim, Attempt와 Lease를 다시 검증해 오래된 Worker의 저장을 차단한다.
+- 최초 실행에서 Version을 `CHUNKED`에서 `EMBEDDING`으로 전환하고 시작 이벤트를 한 번 기록한다.
+- 저장 완료 뒤에도 Version은 `EMBEDDING`을 유지한다.
+- 순차·동시 재호출은 중복 Row를 만들지 않고 기존 전체 결과로 수렴한다.
+
+## 2. 비범위
+
+- Embedding Job의 `INDEXED`, `FAILED` 전환
+- Document Version의 `INDEXED`, `FAILED` 전환
+- Document의 `current_version_id`, 상태 변경
+- Embedding Job Attempt의 `SUCCESS`, `FAILED`, `TIMED_OUT`, `ABANDONED` 전환
+- 검색 가능 Version 교체와 기존 Version Embedding의 `STALE` 전환
+- `EMBEDDING_FAILED`, `INDEXED`, `FAILED`, `RETRY` 이벤트 기록
+- Lease 연장, 만료 Job 회수와 자동 재Claim
+- Batch Embedding API, 병렬 호출, Streaming Insert
+- Chunk별 부분 저장과 Checkpoint 재개
+- 모델별 가변 Vector 컬럼 또는 다중 Dimension 동시 지원
+- Vector 정규화, Quantization과 추가 ANN Index 변경
+- 일반 사용자 API에서 Chunk Text나 Vector를 조회하는 기능
+
+이번 범위는 “모든 Chunk의 Vector Set이 저장돼 다음 인덱싱 완료 단계로 넘어갈 수 있는 상태”까지다.
+검색 가능 상태 확정은 별도 완료 기능이 담당한다.
+
+## 3. 기준선
+
+- GitHub Issue: `#82`
+- Branch: `feature/82`
+- 기준 `develop` Commit: `41014da5be3903013892acda1da95617cea5a53c`
+- Framework: Spring Boot 3.5.16
+- Language: Java 17
+- DB: PostgreSQL/OpenSQL, pgvector, Flyway
+- 기본 Model: `BAAI/bge-m3`, Dimension `1024`
+
+선행 기능은 다음 계약을 제공한다.
+
+- Job Claim은 Worker, Claim Token과 Lease를 `EmbeddingJob`에 기록한다.
+- Attempt 시작은 Job·Worker·Claim Token을 `STARTED` Attempt로 연결한다.
+- Chunk 저장은 Version에 0부터 연속된 불변 Chunk Set을 만들고 상태를 `CHUNKED`로 전환한다.
+- 후속 단계는 Job을 먼저 잠그고 현재 Claim과 Attempt를 검증한다.
+- `embeddings.vector`는 `vector(1024)`이며 `(chunk_id, embedding_model_id)`가 유일하다.
+
+## 4. 핵심 결정
+
+### 4.1 Job 고정 Model 사용
+
+문서 Embedding은 실행 시점의 활성 Model을 다시 조회하지 않는다.
+
+`EmbeddingJob`은 업로드 접수 시점에 사용할 Model을 고정한다. 처리 도중 활성 Model이 교체돼도 이미
+접수된 Job의 Chunk와 검색 Query가 다른 Vector 공간에 들어가지 않도록 Job의 Model ID와 Dimension을
+사용한다.
+
+Job에 Model 연관이나 ID가 없거나 Dimension이 양수가 아니면 사용자 입력 오류가 아니라 내부 설정
+모순으로 처리한다.
+
+### 4.2 API 경로
+
+Endpoint는 다음 경로를 사용한다.
+
+```text
+POST /admin/indexing-jobs/{jobId}/attempts/{attemptId}/embeddings
+```
+
+Attempt ID를 Path에 포함해 현재 Claim 세대의 실행이라는 점을 명시한다. 기존 애플리케이션은 `/api`
+Prefix 없이 `/admin/**`를 사용하며, 기존 ADMIN Security 정책을 그대로 적용한다.
+
+### 4.3 준비·외부 호출·완료 분리
+
+하나의 긴 Transaction 안에서 외부 Embedding Server를 호출하지 않는다.
+
+```text
+준비 Transaction
+ → Transaction 밖 단건 순차 HTTP 호출
+ → 완료 Transaction
+```
+
+준비와 완료 Transaction 모두 Job→Version 순서로 잠근다. 두 Transaction 사이에는 ID, Text, Hash,
+Dimension만 가진 불변 Snapshot과 Draft만 전달하며 JPA Entity를 전달하지 않는다.
+
+### 4.4 전체 Set 단위 저장
+
+MVP는 한 Version의 모든 Chunk Vector를 메모리에서 생성한 후 한 번에 저장한다.
+
+- 중간 외부 호출이 실패하면 완료 Transaction을 호출하지 않는다.
+- 완료 저장 중 한 Row라도 실패하면 전체 Insert를 Rollback한다.
+- 부분 저장 상태는 정상 재개 지점이 아니라 내부 데이터 모순이다.
+- Chunk 수만큼 전부 저장됐을 때만 완료 재생으로 인정한다.
+
+대용량 문서의 Batch·Checkpoint는 별도 저장 상태와 재개 정책이 필요하므로 후속 범위로 분리한다.
+
+### 4.5 단건 순차 외부 호출
+
+Chunk는 `chunk_index` 오름차순으로 한 건씩 호출한다.
+
+- 먼저 실패한 Chunk 이후의 호출은 실행하지 않는다.
+- 호출 순서와 저장 순서가 Chunk Index와 일치한다.
+- Text와 Vector를 로그에 남기지 않는다.
+- 외부 전송 오류는 `EMBEDDING_SERVER_UNAVAILABLE`로 변환한다.
+
+현재 외부 계약은 `POST /embed` 단건 요청이다. Batch와 병렬 처리 없이 가장 작은 재현 가능한 실행
+단위를 유지한다.
+
+### 4.6 상태 경계
+
+정상 상태 변화는 Version에만 적용한다.
+
+| 대상 | 처리 전 | 준비 후 | Vector 저장 후 |
+| --- | --- | --- | --- |
+| `EmbeddingJob` | `PROCESSING` | `PROCESSING` | `PROCESSING` |
+| `EmbeddingJobAttempt` | `STARTED` | `STARTED` | `STARTED` |
+| `DocumentVersion` | `CHUNKED` | `EMBEDDING` | `EMBEDDING` |
+| `Document` | 기존 상태 | 변경 없음 | 변경 없음 |
+
+Version을 `INDEXED`로 바꾸지 않는 이유는 Vector 저장과 검색 가능 Version 교체가 서로 다른 원자성
+경계를 갖기 때문이다. 다음 단계는 Embedding 전체 존재를 다시 검증한 뒤 Job, Attempt, Version,
+Document와 검색 가시성을 함께 확정해야 한다.
+
+### 4.7 Vector Hash
+
+Vector Hash는 다음 규칙으로 계산한다.
+
+1. Vector의 각 `float`를 순서대로 IEEE 754 32-bit 값으로 취급한다.
+2. 각 값을 big-endian 4 Byte로 직렬화한다.
+3. 전체 Byte 배열에 SHA-256을 적용한다.
+4. 소문자 64자리 Hex 문자열로 저장한다.
+
+외부 호출 직후와 완료 저장 직전에 같은 규칙을 사용한다. Draft의 Vector나 Hash가 단계 사이에
+변경되면 저장을 거부한다.
+
+## 5. 불변식
+
+### 5.1 실행 소유권
+
+준비와 완료 Transaction은 다음 조건을 모두 확인한다.
+
+1. Job이 존재하고 `PROCESSING` 상태다.
+2. Job의 현재 Worker, Claim Token, Lock 시작과 Lease 만료 값이 존재한다.
+3. 요청 Worker와 Claim Token이 현재 Job 소유권과 같다.
+4. 잠금 획득 후 현재 시각이 Lease 만료 시각보다 이르다.
+5. Job과 Claim Token으로 조회한 Attempt가 존재한다.
+6. Attempt ID가 Path의 Attempt ID와 같다.
+7. Attempt Worker가 요청 Worker와 같다.
+8. Attempt 상태가 `STARTED`다.
+
+외부 호출 중 Lease가 만료되거나 Claim 세대가 바뀌면 완료 Transaction에서 저장을 거부한다.
+
+### 5.2 Lock 순서
+
+모든 상태 변경 경로는 다음 순서를 유지한다.
+
+1. Embedding Job 쓰기 행 잠금
+2. 잠금 후 현재 시각 계산
+3. Job 소유권과 Lease 검증
+4. 현재 Claim Attempt 검증
+5. Job이 참조하는 Document Version 쓰기 행 잠금
+6. Chunk와 Embedding 저장 상태 검증
+
+Attempt나 Version을 먼저 잠근 뒤 Job을 잠그는 반대 순서를 만들지 않는다.
+
+### 5.3 Chunk Set
+
+- Chunk 목록은 비어 있을 수 없다.
+- Chunk ID와 Version ID가 존재한다.
+- 모든 Chunk는 Job Version을 참조한다.
+- `chunk_index`는 0부터 끊김 없이 증가한다.
+- `chunk_text`는 비어 있지 않다.
+- `content_hash`는 소문자 SHA-256 64자리다.
+- 준비 Snapshot과 완료 시점의 ID, Index, Text, Hash가 모두 같다.
+
+Chunk는 선행 단계에서 확정된 불변 데이터다. 완료 시 달라졌다면 새 데이터를 조용히 사용하지 않고
+내부 모순으로 처리한다.
+
+### 5.4 Vector Draft
+
+- Draft 수는 Chunk 수와 정확히 같다.
+- Draft 순서는 Chunk Index 순서와 같다.
+- Draft의 Chunk ID, Index, Content Hash는 현재 Chunk와 같다.
+- Vector는 null이 아니고 Job Model Dimension과 같다.
+- 모든 원소는 유한 값이며 `NaN`, 양·음의 무한대를 포함하지 않는다.
+- `vector_hash`는 소문자 64자리 Hex다.
+- 저장 직전 다시 계산한 Hash가 Draft Hash와 같다.
+
+Vector 배열은 Generator, Draft와 Entity 경계에서 복사해 호출자가 보관한 배열 변경이 저장 값에
+전파되지 않게 한다.
+
+### 5.5 영속 데이터
+
+각 `Embedding`은 다음 관계를 모두 가진다.
+
+- `chunk`: Vector 원본 Chunk
+- `documentVersion`: Job 대상 Version
+- `document`: Version이 속한 Document
+- `embeddingModel`: Job에 고정된 Model
+- `dimension`: Model Dimension
+- `status`: `ACTIVE`
+- `vectorHash`: 검증된 Vector Hash
+
+역정규화한 `document_id`, `document_version_id`는 Chunk 관계에서 도출한 값과 같아야 한다.
+`(chunk_id, embedding_model_id)` Unique 제약은 Application 잠금 외의 최종 중복 방어선이다.
+
+## 6. 실행 상태 판정
+
+Version 상태와 Job Model의 저장 개수를 함께 판단한다.
+
+| Version 상태 | 현재 Model Embedding 수 | 처리 |
+| --- | ---: | --- |
+| `CHUNKED` | 0 | 최초 작업 시작, `EMBEDDING` 전이와 시작 이벤트 |
+| `CHUNKED` | 1 이상 | 상태·데이터 모순 |
+| `EMBEDDING` | 0 | 외부 호출 실패 또는 중단 이후 처음부터 재개 |
+| `EMBEDDING` | Chunk 수와 같음 | 완료 재생 |
+| `EMBEDDING` | 0과 Chunk 수 사이 | 부분 저장 모순 |
+| 다른 상태 | 0 | 현재 단계 실행 불가 |
+| 다른 상태 | 1 이상 | 상태·데이터 모순 |
+
+`EMBEDDING` 상태의 0개 재개는 첫 Transaction Commit 뒤 Process가 종료되거나 외부 서버 오류가 발생한
+경우를 복구한다. 부분 저장은 이 설계에서 발생할 수 없으므로 별도 오류로 드러낸다.
+
+## 7. 전체 흐름
+
+### 7.1 최초 정상 요청
+
+1. Controller가 양수 Job·Attempt ID와 Worker ID, canonical Claim Token을 검증한다.
+2. 비 Transaction `DocumentEmbeddingService`가 준비 Transaction을 호출한다.
+3. 준비 Transaction이 Job을 잠그고 현재 Claim, Lease와 Attempt를 검증한다.
+4. Job Model을 고정하고 Job Version을 잠근다.
+5. Chunk를 Index 순서로 읽고 전체 Set을 검증한다.
+6. 현재 Model Embedding 수가 0인지 확인한다.
+7. Version을 `CHUNKED`에서 `EMBEDDING`으로 전환한다.
+8. 같은 Transaction에 `EMBEDDING_STARTED` 이벤트를 한 번 저장한다.
+9. Version ID, Model ID, Dimension과 Chunk Snapshot을 반환하고 Commit한다.
+10. DB 행 잠금이 해제된 뒤 Generator가 Chunk Text를 순서대로 외부 서버에 전달한다.
+11. 각 Vector의 차원과 유한 값을 검증하고 Hash를 계산해 Draft를 만든다.
+12. 모든 Chunk 호출이 성공하면 완료 Transaction을 호출한다.
+13. 완료 Transaction이 Job을 다시 잠그고 Claim, Lease와 Attempt를 다시 검증한다.
+14. 같은 Model과 Version인지 확인하고 Version을 잠근다.
+15. Chunk Set과 준비 Snapshot이 같은지 다시 확인한다.
+16. 다른 요청의 선행 완료와 부분 저장 여부를 확인한다.
+17. 모든 Draft와 Vector Hash를 다시 검증한다.
+18. 같은 Version·Model의 `ACTIVE` Embedding Set을 `saveAllAndFlush`로 저장한다.
+19. Version은 `EMBEDDING`으로 유지하고 Commit한다.
+20. Controller가 결과를 `201 Created`로 반환한다.
+
+### 7.2 완료 재생
+
+1. 준비 Transaction에서 소유권과 Attempt를 먼저 검증한다.
+2. Version이 `EMBEDDING`이고 현재 Model Embedding 수가 Chunk 수와 같으면 완료 결과를 반환한다.
+3. 외부 Embedding Server와 완료 Transaction은 호출하지 않는다.
+4. 새 Row와 새 이벤트를 만들지 않는다.
+5. Controller는 같은 응답 Body를 `200 OK`로 반환한다.
+
+### 7.3 외부 호출 실패
+
+1. 준비 Transaction의 `EMBEDDING` 전이와 시작 이벤트는 이미 Commit됐다.
+2. 외부 서버 장애, 차원 불일치 또는 비유한 Vector가 발생하면 나머지 호출을 중단한다.
+3. 완료 Transaction은 호출하지 않는다.
+4. Embedding Row는 하나도 저장되지 않는다.
+5. Version은 `EMBEDDING`, Job은 `PROCESSING`, Attempt는 `STARTED`로 남는다.
+6. 현재 또는 새 유효 Attempt가 같은 Endpoint를 호출해 처음부터 재개할 수 있다.
+
+실패 상태와 Attempt 종료는 후속 실패 처리 기능이 담당한다.
+
+### 7.4 동시 중복 요청
+
+1. 두 요청의 준비 Transaction은 Job 잠금으로 차례로 실행된다.
+2. 첫 요청만 `CHUNKED`에서 `EMBEDDING`으로 전환하고 시작 이벤트를 저장한다.
+3. 두 요청 모두 저장 결과가 0개인 동안 Transaction 밖에서 같은 Vector Set을 계산할 수 있다.
+4. 완료 Transaction은 Job과 Version 잠금으로 다시 직렬화된다.
+5. 먼저 진입한 요청이 전체 Embedding Set을 Commit하고 `201 Created`를 반환한다.
+6. 나중 요청은 전체 저장을 확인하고 Insert 없이 재생해 `200 OK`를 반환한다.
+7. DB에는 Embedding Set 하나와 시작 이벤트 하나만 남는다.
+
+## 8. 계층별 책임
+
+### 8.1 `EmbeddingClient`
+
+- `/embed` 단건 HTTP 전송
+- 외부 전송 오류 변환
+- Model 선택과 Dimension 검증은 수행하지 않음
+- Text와 Vector를 로그에 기록하지 않음
+
+Query Embedding은 활성 Model을 선택하고, 문서 Embedding은 Job 고정 Model을 선택한 뒤 같은 Client를
+재사용한다.
+
+### 8.2 `DocumentEmbeddingGenerator`
+
+- 정렬된 Chunk Snapshot 순차 처리
+- Vector 차원·유한 값 검증
+- Vector Hash 생성
+- 완료 Transaction용 불변 Draft 반환
+- DB와 JPA Entity에 의존하지 않음
+
+### 8.3 `DocumentEmbeddingTransactionService`
+
+- Job→Version 잠금 순서
+- Claim, Lease와 Attempt 검증
+- Version·Model·Chunk·기존 Embedding 상태 판정
+- 최초 `EMBEDDING` 전이와 시작 이벤트
+- 완료 시 Snapshot과 Draft 재검증
+- 전체 Embedding Set 원자 저장
+
+### 8.4 `DocumentEmbeddingService`
+
+- 준비→외부 호출→완료 순서 조정
+- 완료 재생 시 외부 호출 생략
+- 외부 실패 시 완료 Transaction 미호출
+- 내부 완료 결과를 API 응답으로 변환
+- 자체 DB Transaction 없음
+
+### 8.5 `IndexingJobAdminController`
+
+- Path와 Request Body Validation
+- Orchestration 호출
+- 최초 저장 `201 Created`, 재생 `200 OK` 선택
+- Claim Token, Chunk Text와 Vector 비노출
+
+## 9. API 계약
+
+### 9.1 요청
+
+```json
+{
+ "workerId": 1,
+ "claimToken": "34c19d16-6ae1-4f6a-a35d-0123456789ab"
+}
+```
+
+- `jobId`, `attemptId`, `workerId`는 양수다.
+- Claim Token은 소문자 canonical UUID 형식이며 최대 36자다.
+
+### 9.2 응답
+
+```json
+{
+ "success": true,
+ "status": 201,
+ "data": {
+ "jobId": 10,
+ "attemptId": 100,
+ "documentVersionId": 5,
+ "embeddingModelId": 7,
+ "chunkCount": 3,
+ "embeddingCount": 3,
+ "versionStatus": "EMBEDDING"
+ }
+}
+```
+
+응답에는 Claim Token, Chunk Text, Vector와 Vector Hash를 포함하지 않는다.
+
+## 10. 오류 계약
+
+| 상황 | HTTP | Error Code |
+| --- | ---: | --- |
+| Job 없음 | 404 | `EMBEDDING-JOB-001` |
+| Job이 `PROCESSING`이 아님 | 409 | `EMBEDDING-JOB-002` |
+| Worker·Claim Token 불일치 | 409 | `EMBEDDING-JOB-003` |
+| Lease 만료 | 409 | `EMBEDDING-JOB-004` |
+| Job 소유권 데이터 모순 | 500 | `EMBEDDING-JOB-005` |
+| Attempt 불일치 | 409 | `EMBEDDING-JOB-006` |
+| Version 상태에서 실행 불가 | 409 | `DOCUMENT-VERSION-006` |
+| Chunk Set 불일치 | 500 | `DOCUMENT-CHUNK-001` |
+| Version·Embedding Set 불일치 | 500 | `DOCUMENT-EMBEDDING-001` |
+| Vector 또는 Hash가 유효하지 않음 | 500 | `DOCUMENT-EMBEDDING-002` |
+| Job Model 없음 | 500 | `EMBEDDING-MODEL-001` |
+| 외부 서버 장애 | 503 | `SEARCH-001` |
+| Vector 차원 불일치 | 500 | `SEARCH-002` |
+
+## 11. 보안과 운영 주의사항
+
+- Endpoint는 기존 `/admin/**` ADMIN 권한 정책을 사용한다.
+- Claim Token은 요청 소유권 검증에만 사용하고 응답, 이벤트와 로그에 남기지 않는다.
+- Chunk Text와 Vector는 로그, 오류 메시지와 Metric Label에 넣지 않는다.
+- 외부 오류 로그에는 예외 종류만 기록한다.
+- 새 Secret이나 Model 설정을 `application.yml`에 하드코딩하지 않는다.
+- Entity를 Controller 응답으로 직접 노출하지 않는다.
+- `EMBEDDING` 상태와 전체 Row 존재만으로 검색 가능하다고 판단하지 않는다.
+
+## 12. 후속 확장 조건
+
+다음 기능은 별도 설계와 상태 계약이 필요하다.
+
+- 인덱싱 완료 Transaction과 `INDEXED` 전환
+- 새 Version 활성화와 이전 Embedding `STALE` 처리
+- Job·Attempt 성공 및 실패 종료
+- Lease 연장과 만료 Job Recovery
+- Batch Embedding API와 호출 병렬화
+- Chunk 단위 Checkpoint, 부분 재시도와 대용량 메모리 제한
+- 다중 Model Dimension을 위한 물리 Vector 저장 전략
+- Vector 생성 처리량·Latency Metric과 실패율 Dashboard
+
+Batch나 Checkpoint를 도입할 때는 현재의 “부분 저장은 모순” 계약을 그대로 유지할 수 없다. 저장
+상태, 재개 Cursor, 중복 방지 Key와 실패 복구 정책을 먼저 정의해야 한다.
diff --git a/src/main/java/com/opensource/docgrid/domain/document/entity/DocumentVersion.java b/src/main/java/com/opensource/docgrid/domain/document/entity/DocumentVersion.java
index 0576843..e416c5d 100644
--- a/src/main/java/com/opensource/docgrid/domain/document/entity/DocumentVersion.java
+++ b/src/main/java/com/opensource/docgrid/domain/document/entity/DocumentVersion.java
@@ -135,6 +135,10 @@ public void markChunked() {
}
public void markEmbedding() {
+ // Chunk Set이 확정된 Version만 Embedding 생성 단계에 진입할 수 있다.
+ if (status != DocumentVersionStatus.CHUNKED) {
+ throw new IllegalStateException("CHUNKED 상태의 문서 버전만 EMBEDDING으로 전환할 수 있습니다.");
+ }
this.status = DocumentVersionStatus.EMBEDDING;
}
diff --git a/src/main/java/com/opensource/docgrid/domain/document/repository/DocumentChunkRepository.java b/src/main/java/com/opensource/docgrid/domain/document/repository/DocumentChunkRepository.java
index 5884d0f..10158f9 100644
--- a/src/main/java/com/opensource/docgrid/domain/document/repository/DocumentChunkRepository.java
+++ b/src/main/java/com/opensource/docgrid/domain/document/repository/DocumentChunkRepository.java
@@ -1,18 +1,23 @@
package com.opensource.docgrid.domain.document.repository;
+import java.util.List;
+
import org.springframework.data.jpa.repository.JpaRepository;
import com.opensource.docgrid.domain.document.entity.DocumentChunk;
/**
- * 문서 Version별 Chunk 존재 여부와 개수 조회 및 Chunk Set 전체 저장을 담당한다.
+ * 문서 Version별 Chunk 존재 여부·개수·정렬 조회와 Chunk Set 전체 저장을 담당한다.
*
*
Chunk 생성 Transaction은 기존 결과 확인에 존재·개수 조회를 사용하고, 신규 결과는
* {@link JpaRepository#saveAllAndFlush(Iterable)}로 같은 Transaction 안에서 즉시 검증한다.
+ * Embedding 생성은 Chunk 순서를 저장 결과와 일치시키기 위해 chunkIndex 오름차순 조회를 사용한다.
*/
public interface DocumentChunkRepository extends JpaRepository {
boolean existsByDocumentVersionId(Long documentVersionId);
long countByDocumentVersionId(Long documentVersionId);
+
+ List findAllByDocumentVersionIdOrderByChunkIndexAsc(Long documentVersionId);
}
diff --git a/src/main/java/com/opensource/docgrid/domain/embedding/client/EmbeddingClient.java b/src/main/java/com/opensource/docgrid/domain/embedding/client/EmbeddingClient.java
new file mode 100644
index 0000000..7aac3d5
--- /dev/null
+++ b/src/main/java/com/opensource/docgrid/domain/embedding/client/EmbeddingClient.java
@@ -0,0 +1,50 @@
+package com.opensource.docgrid.domain.embedding.client;
+
+import org.springframework.beans.factory.annotation.Qualifier;
+import org.springframework.stereotype.Component;
+import org.springframework.web.client.RestClient;
+import org.springframework.web.client.RestClientException;
+
+import com.opensource.docgrid.domain.embedding.dto.request.EmbedRequest;
+import com.opensource.docgrid.domain.embedding.dto.response.EmbedServerResponse;
+import com.opensource.docgrid.global.exception.DocGridException;
+import com.opensource.docgrid.global.exception.ErrorCode;
+
+import lombok.extern.slf4j.Slf4j;
+
+/**
+ * 외부 Embedding Server의 단건 Vector 생성 HTTP 계약을 담당한다.
+ *
+ * 이 Client는 모델 선택과 Vector 차원 검증을 수행하지 않는다. 호출 Service가 실행 Context에 맞는
+ * 모델을 선택하고 반환 Vector를 검증하며, 이 클래스는 전송 오류를 공통 서비스 장애로 변환하는 경계만
+ * 책임진다.
+ */
+@Slf4j
+@Component
+public class EmbeddingClient {
+
+ private final RestClient restClient;
+
+ public EmbeddingClient(@Qualifier("embeddingRestClient") RestClient restClient) {
+ this.restClient = restClient;
+ }
+
+ /**
+ * 입력 Text를 외부 서버에 전달하고 Dense Vector를 반환한다.
+ */
+ public float[] embed(String text) {
+ EmbedServerResponse response;
+ try {
+ response = restClient.post()
+ .uri("/embed")
+ .body(new EmbedRequest(text))
+ .retrieve()
+ .body(EmbedServerResponse.class);
+ } catch (RestClientException exception) {
+ log.error("임베딩 서버 호출에 실패했습니다. cause={}", exception.getClass().getSimpleName());
+ throw new DocGridException(ErrorCode.EMBEDDING_SERVER_UNAVAILABLE);
+ }
+
+ return response == null ? null : response.vector();
+ }
+}
diff --git a/src/main/java/com/opensource/docgrid/domain/embedding/controller/IndexingJobAdminController.java b/src/main/java/com/opensource/docgrid/domain/embedding/controller/IndexingJobAdminController.java
index 0bba777..413bf61 100644
--- a/src/main/java/com/opensource/docgrid/domain/embedding/controller/IndexingJobAdminController.java
+++ b/src/main/java/com/opensource/docgrid/domain/embedding/controller/IndexingJobAdminController.java
@@ -14,11 +14,15 @@
import com.opensource.docgrid.domain.embedding.dto.request.StartEmbeddingJobAttemptRequest;
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.response.ClaimedEmbeddingJobResponse;
import com.opensource.docgrid.domain.embedding.dto.response.DocumentChunksResponse;
+import com.opensource.docgrid.domain.embedding.dto.response.DocumentEmbeddingsResponse;
import com.opensource.docgrid.domain.embedding.dto.response.StartedEmbeddingJobAttemptResponse;
import com.opensource.docgrid.domain.document.service.DocumentParsingService;
import com.opensource.docgrid.domain.document.service.command.DocumentChunkTransactionService.ChunkResult;
+import com.opensource.docgrid.domain.embedding.service.DocumentEmbeddingService;
+import com.opensource.docgrid.domain.embedding.service.DocumentEmbeddingService.EmbeddingResult;
import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobClaimService;
import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobAttemptService;
import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobAttemptService.StartResult;
@@ -36,10 +40,10 @@
import lombok.RequiredArgsConstructor;
/**
- * 관리자용 Embedding Job Claim과 Attempt 시작 요청을 HTTP API로 제공하는 Controller.
+ * 관리자용 Embedding Job Claim, Attempt 시작과 문서 Chunk·Embedding 실행을 HTTP API로 제공한다.
*
- *
HTTP 입력 검증과 성공 상태 변환만 담당한다. Job Claim 및 현재 소유권 기반 Attempt 시작의
- * Transaction·동시성 규칙은 각 Command Service에 위임한다.
+ *
HTTP 입력 검증과 성공 상태 변환만 담당한다. Job Claim 및 현재 소유권 기반 파이프라인 단계의
+ * Transaction·외부 호출·동시성 규칙은 각 Service에 위임한다.
*/
@Tag(name = "Admin - Indexing Job", description = "관리자 전용 인덱싱 Job 제어 API")
@Validated
@@ -51,6 +55,7 @@ public class IndexingJobAdminController {
private final EmbeddingJobClaimService embeddingJobClaimService;
private final EmbeddingJobAttemptService embeddingJobAttemptService;
private final DocumentParsingService documentParsingService;
+ private final DocumentEmbeddingService documentEmbeddingService;
@Operation(
summary = "PENDING Job Claim",
@@ -227,4 +232,70 @@ public ResponseEntity> createChunks(
}
return ResponseUtils.ok(result.response());
}
+
+ @Operation(
+ summary = "Document Chunk Embedding 생성",
+ description = "현재 PROCESSING Job의 유효한 Attempt 소유권과 Job 고정 Model을 검증하고 "
+ + "Chunk를 순서대로 외부 Embedding 서버에 전달한 뒤 Vector Set을 원자 저장합니다. "
+ + "최초 저장은 201, 기존 완료 결과의 멱등 재생은 200을 반환합니다."
+ )
+ @ApiResponses({
+ @io.swagger.v3.oas.annotations.responses.ApiResponse(
+ responseCode = "201",
+ description = "Document Embedding 최초 저장"
+ ),
+ @io.swagger.v3.oas.annotations.responses.ApiResponse(
+ responseCode = "200",
+ description = "기존 Embedding 결과 재생"
+ ),
+ @io.swagger.v3.oas.annotations.responses.ApiResponse(
+ responseCode = "400",
+ description = "Job ID, Attempt ID, Worker ID 또는 Claim Token 형식 오류",
+ content = @Content(schema = @Schema(implementation = ErrorResponse.class))
+ ),
+ @io.swagger.v3.oas.annotations.responses.ApiResponse(
+ responseCode = "403",
+ description = "인증되지 않았거나 ADMIN 권한 없음",
+ content = @Content(schema = @Schema(implementation = ErrorResponse.class))
+ ),
+ @io.swagger.v3.oas.annotations.responses.ApiResponse(
+ responseCode = "404",
+ description = "Embedding Job 없음",
+ content = @Content(schema = @Schema(implementation = ErrorResponse.class))
+ ),
+ @io.swagger.v3.oas.annotations.responses.ApiResponse(
+ responseCode = "409",
+ description = "현재 소유권, Attempt, Lease 또는 Version 상태 오류",
+ content = @Content(schema = @Schema(implementation = ErrorResponse.class))
+ ),
+ @io.swagger.v3.oas.annotations.responses.ApiResponse(
+ responseCode = "500",
+ description = "Model, Chunk, Vector 또는 Embedding 저장 상태 불일치",
+ content = @Content(schema = @Schema(implementation = ErrorResponse.class))
+ ),
+ @io.swagger.v3.oas.annotations.responses.ApiResponse(
+ responseCode = "503",
+ description = "외부 Embedding 서버 장애",
+ content = @Content(schema = @Schema(implementation = ErrorResponse.class))
+ )
+ })
+ @PostMapping(
+ value = "/{jobId}/attempts/{attemptId}/embeddings",
+ consumes = MediaType.APPLICATION_JSON_VALUE,
+ produces = MediaType.APPLICATION_JSON_VALUE
+ )
+ public ResponseEntity> createEmbeddings(
+ @PathVariable @Positive Long jobId,
+ @PathVariable @Positive Long attemptId,
+ @Valid @RequestBody CreateDocumentEmbeddingsRequest request
+ ) {
+ // 1. 비 Transaction Service가 준비·외부 호출·완료 Transaction의 순서를 조정한다.
+ EmbeddingResult result = documentEmbeddingService.createEmbeddings(jobId, attemptId, request);
+
+ // 2. 같은 응답 Body를 사용하고 실제 최초 저장 여부로 HTTP 상태만 구분한다.
+ if (result.created()) {
+ return ResponseUtils.created(result.response());
+ }
+ return ResponseUtils.ok(result.response());
+ }
}
diff --git a/src/main/java/com/opensource/docgrid/domain/embedding/dto/request/CreateDocumentEmbeddingsRequest.java b/src/main/java/com/opensource/docgrid/domain/embedding/dto/request/CreateDocumentEmbeddingsRequest.java
new file mode 100644
index 0000000..14c72fe
--- /dev/null
+++ b/src/main/java/com/opensource/docgrid/domain/embedding/dto/request/CreateDocumentEmbeddingsRequest.java
@@ -0,0 +1,32 @@
+package com.opensource.docgrid.domain.embedding.dto.request;
+
+import io.swagger.v3.oas.annotations.media.Schema;
+import jakarta.validation.constraints.NotBlank;
+import jakarta.validation.constraints.NotNull;
+import jakarta.validation.constraints.Pattern;
+import jakarta.validation.constraints.Positive;
+import jakarta.validation.constraints.Size;
+
+/**
+ * 현재 Embedding Job Attempt 소유권으로 문서 Chunk Embedding 생성을 요청하는 DTO.
+ *
+ * Worker ID와 canonical UUID Claim Token은 준비와 완료 단계의 소유권 검증에만 사용하며,
+ * 응답, 이벤트와 로그에는 노출하지 않는다.
+ */
+public record CreateDocumentEmbeddingsRequest(
+ @Schema(description = "현재 Job을 소유한 Worker 식별자", example = "1")
+ @NotNull
+ @Positive
+ Long workerId,
+
+ @Schema(description = "현재 Claim의 canonical UUID Token",
+ example = "34c19d16-6ae1-4f6a-a35d-0123456789ab")
+ @NotBlank
+ @Size(max = 36)
+ @Pattern(
+ regexp = "^[0-9a-f]{8}-[0-9a-f]{4}-[1-5][0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$",
+ message = "canonical UUID 형식이어야 합니다."
+ )
+ String claimToken
+) {
+}
diff --git a/src/main/java/com/opensource/docgrid/domain/embedding/dto/response/DocumentEmbeddingsResponse.java b/src/main/java/com/opensource/docgrid/domain/embedding/dto/response/DocumentEmbeddingsResponse.java
new file mode 100644
index 0000000..ae3cf46
--- /dev/null
+++ b/src/main/java/com/opensource/docgrid/domain/embedding/dto/response/DocumentEmbeddingsResponse.java
@@ -0,0 +1,35 @@
+package com.opensource.docgrid.domain.embedding.dto.response;
+
+import com.opensource.docgrid.domain.document.enums.DocumentVersionStatus;
+
+import io.swagger.v3.oas.annotations.media.Schema;
+
+/**
+ * 생성됐거나 멱등 재생된 Document Embedding Set의 식별자와 집계 상태를 전달한다.
+ *
+ *
Claim Token, Chunk Text와 Vector는 제외하고 현재 실행 Context, 고정 Model과 Version 결과만
+ * 노출한다. Embedding 저장 완료 시에도 Version은 후속 색인 완료 전까지 EMBEDDING을 유지한다.
+ */
+public record DocumentEmbeddingsResponse(
+ @Schema(description = "처리한 Embedding Job 식별자", example = "10")
+ Long jobId,
+
+ @Schema(description = "현재 실행 Attempt 식별자", example = "100")
+ Long attemptId,
+
+ @Schema(description = "Embedding이 저장된 Document Version 식별자", example = "5")
+ Long documentVersionId,
+
+ @Schema(description = "Job에 고정된 Embedding Model 식별자", example = "7")
+ Long embeddingModelId,
+
+ @Schema(description = "Version의 전체 Chunk 수", example = "3")
+ int chunkCount,
+
+ @Schema(description = "현재 Model로 저장된 Embedding 수", example = "3")
+ int embeddingCount,
+
+ @Schema(description = "Embedding 저장 후 Version 상태", example = "EMBEDDING")
+ DocumentVersionStatus versionStatus
+) {
+}
diff --git a/src/main/java/com/opensource/docgrid/domain/embedding/entity/Embedding.java b/src/main/java/com/opensource/docgrid/domain/embedding/entity/Embedding.java
index 56df2ac..a67e327 100644
--- a/src/main/java/com/opensource/docgrid/domain/embedding/entity/Embedding.java
+++ b/src/main/java/com/opensource/docgrid/domain/embedding/entity/Embedding.java
@@ -1,5 +1,7 @@
package com.opensource.docgrid.domain.embedding.entity;
+import java.util.Arrays;
+
import com.opensource.docgrid.domain.document.entity.Document;
import com.opensource.docgrid.domain.document.entity.DocumentChunk;
import com.opensource.docgrid.domain.document.entity.DocumentVersion;
@@ -108,9 +110,21 @@ public Embedding(DocumentChunk chunk, Document document, DocumentVersion documen
this.document = document;
this.documentVersion = documentVersion;
this.embeddingModel = embeddingModel;
- this.vector = vector;
+ // 호출자가 보관한 배열 변경이 영속화 값에 전파되지 않도록 생성 시점에 복사한다.
+ this.vector = copyVector(vector);
this.dimension = dimension;
this.vectorHash = vectorHash;
this.status = status != null ? status : EmbeddingStatus.ACTIVE;
}
+
+ /**
+ * 영속 Entity 내부 Vector가 호출자에 의해 변경되지 않도록 복사본을 반환한다.
+ */
+ public float[] getVector() {
+ return copyVector(vector);
+ }
+
+ private static float[] copyVector(float[] source) {
+ return source == null ? null : Arrays.copyOf(source, source.length);
+ }
}
diff --git a/src/main/java/com/opensource/docgrid/domain/embedding/repository/EmbeddingRepository.java b/src/main/java/com/opensource/docgrid/domain/embedding/repository/EmbeddingRepository.java
new file mode 100644
index 0000000..6d4e4f6
--- /dev/null
+++ b/src/main/java/com/opensource/docgrid/domain/embedding/repository/EmbeddingRepository.java
@@ -0,0 +1,19 @@
+package com.opensource.docgrid.domain.embedding.repository;
+
+import org.springframework.data.jpa.repository.JpaRepository;
+
+import com.opensource.docgrid.domain.embedding.entity.Embedding;
+
+/**
+ * 문서 Chunk Embedding Set의 영속화와 Version·Model 단위 저장 개수 조회를 담당한다.
+ *
+ *
Embedding 생성 Transaction은 Job에 고정된 Model 범위의 저장 개수로 최초 실행, 재개,
+ * 완료 재생과 부분 저장 모순을 구분하고 신규 Set은 한 Transaction에서 전체 저장한다.
+ */
+public interface EmbeddingRepository extends JpaRepository {
+
+ long countByDocumentVersionIdAndEmbeddingModelId(
+ Long documentVersionId,
+ Long embeddingModelId
+ );
+}
diff --git a/src/main/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingDraft.java b/src/main/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingDraft.java
new file mode 100644
index 0000000..58ac4f7
--- /dev/null
+++ b/src/main/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingDraft.java
@@ -0,0 +1,31 @@
+package com.opensource.docgrid.domain.embedding.service;
+
+import java.util.Arrays;
+
+/**
+ * 외부 호출로 생성돼 아직 영속화되지 않은 단일 Chunk Embedding의 불변 값을 전달한다.
+ *
+ * 완료 Transaction은 Chunk ID·순서·내용 Hash를 준비 Snapshot과 다시 비교하고 Vector 차원,
+ * 유한 값과 Hash를 검증한 뒤에만 이 값을 Entity로 변환한다.
+ */
+public record DocumentEmbeddingDraft(
+ Long chunkId,
+ int chunkIndex,
+ String contentHash,
+ float[] vector,
+ String vectorHash
+) {
+
+ public DocumentEmbeddingDraft {
+ vector = copyVector(vector);
+ }
+
+ @Override
+ public float[] vector() {
+ return copyVector(vector);
+ }
+
+ private static float[] copyVector(float[] source) {
+ return source == null ? null : Arrays.copyOf(source, source.length);
+ }
+}
diff --git a/src/main/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingGenerator.java b/src/main/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingGenerator.java
new file mode 100644
index 0000000..0e926e5
--- /dev/null
+++ b/src/main/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingGenerator.java
@@ -0,0 +1,65 @@
+package com.opensource.docgrid.domain.embedding.service;
+
+import java.util.ArrayList;
+import java.util.List;
+
+import org.springframework.stereotype.Service;
+
+import com.opensource.docgrid.domain.embedding.client.EmbeddingClient;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.ChunkSnapshot;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.EmbeddingWork;
+import com.opensource.docgrid.global.exception.DocGridException;
+import com.opensource.docgrid.global.exception.ErrorCode;
+
+import lombok.RequiredArgsConstructor;
+
+/**
+ * 준비 Snapshot의 Chunk를 순서대로 외부 서버에 전달하고 검증된 Vector Draft를 생성한다.
+ *
+ *
DB Transaction과 JPA Entity를 사용하지 않으며, 한 번에 한 Chunk만 호출해 실패 시 어떤 결과도
+ * 저장되지 않게 한다. 반환 Vector는 Model 차원과 유한 값을 검증한 뒤 SHA-256 Hash와 함께 복사한다.
+ */
+@Service
+@RequiredArgsConstructor
+public class DocumentEmbeddingGenerator {
+
+ private final EmbeddingClient embeddingClient;
+
+ /**
+ * 정렬된 Chunk Snapshot을 단건 순차 호출해 같은 순서의 Embedding Draft로 변환한다.
+ */
+ public List generate(EmbeddingWork work) {
+ validateWork(work);
+
+ List drafts = new ArrayList<>(work.chunks().size());
+ for (ChunkSnapshot chunk : work.chunks()) {
+ // 1. 현재 Chunk Text만 외부 서버로 보내 DB Transaction 없이 Vector를 생성한다.
+ float[] vector = embeddingClient.embed(chunk.chunkText());
+
+ // 2. Job 고정 Model의 차원과 모든 원소의 유한성을 저장 전에 검증한다.
+ EmbeddingVectorSupport.validate(vector, work.dimension());
+
+ // 3. 검증된 Vector와 원본 Chunk Snapshot을 결합해 완료 Transaction용 Draft를 만든다.
+ drafts.add(new DocumentEmbeddingDraft(
+ chunk.chunkId(),
+ chunk.chunkIndex(),
+ chunk.contentHash(),
+ vector,
+ EmbeddingVectorSupport.calculateHash(vector)
+ ));
+ }
+ return List.copyOf(drafts);
+ }
+
+ private void validateWork(EmbeddingWork work) {
+ if (work == null
+ || work.documentVersionId() == null
+ || work.embeddingModelId() == null
+ || work.dimension() <= 0
+ || work.chunks() == null
+ || work.chunks().isEmpty()) {
+ throw new DocGridException(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT);
+ }
+ }
+
+}
diff --git a/src/main/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingService.java b/src/main/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingService.java
new file mode 100644
index 0000000..f0a5e38
--- /dev/null
+++ b/src/main/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingService.java
@@ -0,0 +1,86 @@
+package com.opensource.docgrid.domain.embedding.service;
+
+import java.util.List;
+
+import org.springframework.stereotype.Service;
+
+import com.opensource.docgrid.domain.embedding.dto.request.CreateDocumentEmbeddingsRequest;
+import com.opensource.docgrid.domain.embedding.dto.response.DocumentEmbeddingsResponse;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.CompletionResult;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.PreparationResult;
+
+import lombok.RequiredArgsConstructor;
+
+/**
+ * 두 개의 짧은 DB Transaction 사이에서 Chunk별 외부 Embedding 호출을 조정한다.
+ *
+ * 이 Service 자체에는 Transaction을 적용하지 않아 외부 HTTP 호출 중 DB 행 잠금이 유지되지 않게
+ * 한다. 준비 Snapshot과 Vector Draft만 단계 사이에 전달하고 JPA Entity는 전달하지 않는다.
+ */
+@Service
+@RequiredArgsConstructor
+public class DocumentEmbeddingService {
+
+ private final DocumentEmbeddingTransactionService transactionService;
+ private final DocumentEmbeddingGenerator embeddingGenerator;
+
+ /**
+ * 현재 Attempt가 소유한 Job의 Chunk Embedding을 생성·저장하거나 기존 결과를 재생한다.
+ */
+ public EmbeddingResult createEmbeddings(
+ Long jobId,
+ Long attemptId,
+ CreateDocumentEmbeddingsRequest request
+ ) {
+ // 1. 준비 Transaction에서 소유권, Version·Model과 Chunk Set을 검증한다.
+ PreparationResult preparation = transactionService.prepare(
+ jobId,
+ attemptId,
+ request.workerId(),
+ request.claimToken()
+ );
+
+ // 2. 같은 Model의 전체 결과가 이미 저장됐으면 외부 서버를 호출하지 않고 즉시 재생한다.
+ if (preparation.isReplay()) {
+ return result(preparation.replayResult());
+ }
+
+ // 3. Transaction 밖에서 Chunk를 순서대로 호출해 검증된 Vector Draft 전체를 만든다.
+ List drafts = embeddingGenerator.generate(preparation.work());
+
+ // 4. 완료 Transaction이 소유권과 Snapshot을 다시 검증하고 전체 Set을 원자 저장한다.
+ return result(transactionService.complete(
+ jobId,
+ attemptId,
+ request.workerId(),
+ request.claimToken(),
+ preparation.work(),
+ drafts
+ ));
+ }
+
+ private EmbeddingResult result(CompletionResult completion) {
+ return new EmbeddingResult(
+ new DocumentEmbeddingsResponse(
+ completion.jobId(),
+ completion.attemptId(),
+ completion.documentVersionId(),
+ completion.embeddingModelId(),
+ completion.chunkCount(),
+ completion.embeddingCount(),
+ completion.documentVersionStatus()
+ ),
+ completion.created()
+ );
+ }
+
+ /**
+ * Embedding Set 응답과 HTTP 생성·재생 상태를 Controller에 함께 전달한다.
+ */
+ public record EmbeddingResult(
+ DocumentEmbeddingsResponse response,
+ boolean created
+ ) {
+ }
+}
diff --git a/src/main/java/com/opensource/docgrid/domain/embedding/service/EmbeddingVectorSupport.java b/src/main/java/com/opensource/docgrid/domain/embedding/service/EmbeddingVectorSupport.java
new file mode 100644
index 0000000..6470f80
--- /dev/null
+++ b/src/main/java/com/opensource/docgrid/domain/embedding/service/EmbeddingVectorSupport.java
@@ -0,0 +1,54 @@
+package com.opensource.docgrid.domain.embedding.service;
+
+import java.nio.ByteBuffer;
+import java.nio.ByteOrder;
+import java.security.MessageDigest;
+import java.security.NoSuchAlgorithmException;
+import java.util.HexFormat;
+
+import com.opensource.docgrid.global.exception.DocGridException;
+import com.opensource.docgrid.global.exception.ErrorCode;
+
+/**
+ * 문서 Embedding 생성과 저장 양쪽에서 사용하는 Vector 차원·유한 값 검증과 Hash 규칙을 제공한다.
+ *
+ * Hash는 float 원소를 순서대로 big-endian IEEE 754 byte로 직렬화한 뒤 SHA-256을 적용한다.
+ * 외부 호출 직후와 영속화 직전에 같은 규칙을 실행해 변경되거나 손상된 Draft 저장을 차단한다.
+ */
+public final class EmbeddingVectorSupport {
+
+ private EmbeddingVectorSupport() {
+ }
+
+ /**
+ * Vector가 Model 차원과 일치하고 모든 원소가 유한한지 검증한다.
+ */
+ public static void validate(float[] vector, int expectedDimension) {
+ if (vector == null || vector.length != expectedDimension) {
+ throw new DocGridException(ErrorCode.EMBEDDING_DIMENSION_MISMATCH);
+ }
+ for (float value : vector) {
+ if (!Float.isFinite(value)) {
+ throw new DocGridException(ErrorCode.EMBEDDING_VECTOR_INVALID);
+ }
+ }
+ }
+
+ /**
+ * 검증된 Vector의 재현 가능한 SHA-256 Hash를 계산한다.
+ */
+ public static String calculateHash(float[] vector) {
+ try {
+ MessageDigest digest = MessageDigest.getInstance("SHA-256");
+ ByteBuffer buffer = ByteBuffer
+ .allocate(vector.length * Float.BYTES)
+ .order(ByteOrder.BIG_ENDIAN);
+ for (float value : vector) {
+ buffer.putFloat(value);
+ }
+ return HexFormat.of().formatHex(digest.digest(buffer.array()));
+ } catch (NoSuchAlgorithmException exception) {
+ throw new DocGridException(ErrorCode.EMBEDDING_VECTOR_INVALID, exception);
+ }
+ }
+}
diff --git a/src/main/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionService.java b/src/main/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionService.java
new file mode 100644
index 0000000..a925562
--- /dev/null
+++ b/src/main/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionService.java
@@ -0,0 +1,443 @@
+package com.opensource.docgrid.domain.embedding.service.command;
+
+import java.time.Clock;
+import java.time.LocalDateTime;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Objects;
+
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+import org.springframework.util.StringUtils;
+
+import com.opensource.docgrid.domain.document.entity.DocumentChunk;
+import com.opensource.docgrid.domain.document.entity.DocumentVersion;
+import com.opensource.docgrid.domain.document.enums.DocumentVersionStatus;
+import com.opensource.docgrid.domain.document.repository.DocumentChunkRepository;
+import com.opensource.docgrid.domain.document.repository.DocumentVersionRepository;
+import com.opensource.docgrid.domain.embedding.entity.Embedding;
+import com.opensource.docgrid.domain.embedding.entity.EmbeddingJob;
+import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel;
+import com.opensource.docgrid.domain.embedding.enums.EmbeddingStatus;
+import com.opensource.docgrid.domain.embedding.repository.EmbeddingJobRepository;
+import com.opensource.docgrid.domain.embedding.repository.EmbeddingRepository;
+import com.opensource.docgrid.domain.embedding.service.DocumentEmbeddingDraft;
+import com.opensource.docgrid.domain.embedding.service.EmbeddingVectorSupport;
+import com.opensource.docgrid.domain.worker.entity.EmbeddingJobAttempt;
+import com.opensource.docgrid.domain.worker.entity.IndexingEvent;
+import com.opensource.docgrid.domain.worker.enums.AttemptStatus;
+import com.opensource.docgrid.domain.worker.enums.IndexingEventType;
+import com.opensource.docgrid.domain.worker.repository.EmbeddingJobAttemptRepository;
+import com.opensource.docgrid.domain.worker.repository.IndexingEventRepository;
+import com.opensource.docgrid.global.exception.DocGridException;
+import com.opensource.docgrid.global.exception.ErrorCode;
+
+import lombok.RequiredArgsConstructor;
+
+/**
+ * 문서 Chunk Embedding 생성의 준비와 완료 단계를 짧은 DB Transaction으로 분리해 수행한다.
+ *
+ *
두 단계는 Job을 먼저, Version을 다음 순서로 잠가 Claim 교체와 같은 Version의 동시 실행을
+ * 직렬화한다. 준비 단계는 Job에 고정된 Model과 Chunk 불변 Snapshot만 외부 호출 구간에 전달한다.
+ */
+@Service
+@RequiredArgsConstructor
+@Transactional
+public class DocumentEmbeddingTransactionService {
+
+ private static final String EMBEDDING_STARTED_MESSAGE = "Document Version Embedding 생성을 시작했습니다.";
+
+ private final EmbeddingJobRepository embeddingJobRepository;
+ private final EmbeddingJobAttemptRepository embeddingJobAttemptRepository;
+ private final DocumentVersionRepository documentVersionRepository;
+ private final DocumentChunkRepository documentChunkRepository;
+ private final EmbeddingRepository embeddingRepository;
+ private final IndexingEventRepository indexingEventRepository;
+ private final EmbeddingJobOwnershipValidator ownershipValidator;
+ private final Clock clock;
+
+ /**
+ * 외부 Embedding 호출 전에 실행 소유권과 저장 상태를 검증하고 불변 작업 Snapshot을 만든다.
+ */
+ public PreparationResult prepare(
+ Long jobId,
+ Long attemptId,
+ Long workerId,
+ String claimToken
+ ) {
+ // 1. Claim 교체와 같은 Job의 중복 요청을 직렬화하도록 Job을 먼저 잠근다.
+ EmbeddingJob embeddingJob = findLockedJob(jobId);
+ LocalDateTime preparedAt = LocalDateTime.now(clock);
+ ownershipValidator.validate(embeddingJob, workerId, claimToken, preparedAt);
+ validateAttempt(embeddingJob, attemptId, workerId, claimToken);
+
+ // 2. Job이 직접 가리키는 Version과 Model을 고정하고 Version 행을 잠근다.
+ EmbeddingModel embeddingModel = findJobModel(embeddingJob);
+ DocumentVersion documentVersion = findLockedJobVersion(embeddingJob);
+
+ // 3. Chunk Set 전체를 순서대로 검증해 부분·중복·누락된 입력을 외부 호출 전에 차단한다.
+ List chunks = documentChunkRepository
+ .findAllByDocumentVersionIdOrderByChunkIndexAsc(documentVersion.getId());
+ validateChunks(documentVersion, chunks);
+
+ // 4. Version 상태와 현재 Model의 저장 개수를 함께 비교해 작업·재개·완료 재생을 구분한다.
+ EmbeddingState state = resolveState(documentVersion, embeddingModel, chunks.size());
+ if (state == EmbeddingState.REPLAY) {
+ return PreparationResult.replay(result(
+ jobId,
+ attemptId,
+ documentVersion,
+ embeddingModel,
+ chunks.size(),
+ false
+ ));
+ }
+
+ // 5. 최초 CHUNKED 요청만 EMBEDDING 상태와 시작 이벤트를 같은 Transaction에 기록한다.
+ if (documentVersion.getStatus() == DocumentVersionStatus.CHUNKED) {
+ documentVersion.markEmbedding();
+ indexingEventRepository.save(IndexingEvent.builder()
+ .embeddingJob(embeddingJob)
+ .eventType(IndexingEventType.EMBEDDING_STARTED)
+ .fromStatus(DocumentVersionStatus.CHUNKED.name())
+ .toStatus(DocumentVersionStatus.EMBEDDING.name())
+ .message(EMBEDDING_STARTED_MESSAGE)
+ .occurredAt(preparedAt)
+ .build());
+ }
+
+ return PreparationResult.work(new EmbeddingWork(
+ documentVersion.getId(),
+ embeddingModel.getId(),
+ embeddingModel.getDimension(),
+ chunks.stream()
+ .map(chunk -> new ChunkSnapshot(
+ chunk.getId(),
+ chunk.getChunkIndex(),
+ chunk.getChunkText(),
+ chunk.getContentHash()
+ ))
+ .toList()
+ ));
+ }
+
+ /**
+ * 외부 호출 결과를 현재 실행 소유권과 준비 Snapshot으로 재검증한 뒤 한 번에 저장한다.
+ */
+ public CompletionResult complete(
+ Long jobId,
+ Long attemptId,
+ Long workerId,
+ String claimToken,
+ EmbeddingWork preparedWork,
+ List drafts
+ ) {
+ // 1. 외부 호출 중 Claim 교체나 Lease 만료를 차단하도록 Job과 Attempt를 다시 검증한다.
+ EmbeddingJob embeddingJob = findLockedJob(jobId);
+ LocalDateTime completedAt = LocalDateTime.now(clock);
+ ownershipValidator.validate(embeddingJob, workerId, claimToken, completedAt);
+ validateAttempt(embeddingJob, attemptId, workerId, claimToken);
+
+ // 2. Job 고정 Model과 Version을 다시 확인하고 준비 단계와 같은 대상인지 검증한다.
+ EmbeddingModel embeddingModel = findJobModel(embeddingJob);
+ DocumentVersion documentVersion = findLockedJobVersion(embeddingJob);
+ validatePreparedTarget(preparedWork, documentVersion, embeddingModel);
+
+ // 3. 현재 Chunk Set을 다시 읽어 외부 호출 중 원본이 바뀌거나 누락되지 않았는지 확인한다.
+ List chunks = documentChunkRepository
+ .findAllByDocumentVersionIdOrderByChunkIndexAsc(documentVersion.getId());
+ validateChunks(documentVersion, chunks);
+ validatePreparedChunks(preparedWork, chunks);
+
+ // 4. 동시 요청이 먼저 전체 저장했으면 기존 결과를 재생하고 부분 저장은 내부 모순으로 거부한다.
+ EmbeddingState state = resolveState(documentVersion, embeddingModel, chunks.size());
+ if (state == EmbeddingState.REPLAY) {
+ return result(jobId, attemptId, documentVersion, embeddingModel, chunks.size(), false);
+ }
+ if (documentVersion.getStatus() != DocumentVersionStatus.EMBEDDING) {
+ throw new DocGridException(ErrorCode.DOCUMENT_VERSION_EMBEDDING_NOT_ALLOWED);
+ }
+
+ // 5. 모든 Draft를 다시 검증하고 같은 Version·Model의 ACTIVE Embedding Set으로 원자 저장한다.
+ List embeddings = toEntities(
+ documentVersion,
+ embeddingModel,
+ chunks,
+ drafts
+ );
+ embeddingRepository.saveAllAndFlush(embeddings);
+
+ return result(jobId, attemptId, documentVersion, embeddingModel, embeddings.size(), true);
+ }
+
+ private EmbeddingJob findLockedJob(Long jobId) {
+ return embeddingJobRepository.findByIdForUpdate(jobId)
+ .orElseThrow(() -> new DocGridException(ErrorCode.EMBEDDING_JOB_NOT_FOUND));
+ }
+
+ private DocumentVersion findLockedJobVersion(EmbeddingJob embeddingJob) {
+ if (embeddingJob.getDocumentVersion() == null
+ || embeddingJob.getDocumentVersion().getId() == null) {
+ throw new DocGridException(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT);
+ }
+ return documentVersionRepository.findByIdForUpdate(embeddingJob.getDocumentVersion().getId())
+ .orElseThrow(() -> new DocGridException(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT));
+ }
+
+ private EmbeddingModel findJobModel(EmbeddingJob embeddingJob) {
+ EmbeddingModel embeddingModel = embeddingJob.getEmbeddingModel();
+ if (embeddingModel == null
+ || embeddingModel.getId() == null
+ || embeddingModel.getDimension() <= 0) {
+ throw new DocGridException(ErrorCode.EMBEDDING_MODEL_NOT_CONFIGURED);
+ }
+ return embeddingModel;
+ }
+
+ private void validateAttempt(
+ EmbeddingJob embeddingJob,
+ Long attemptId,
+ Long workerId,
+ String claimToken
+ ) {
+ EmbeddingJobAttempt attempt = embeddingJobAttemptRepository
+ .findByEmbeddingJobIdAndClaimToken(embeddingJob.getId(), claimToken)
+ .orElseThrow(() -> new DocGridException(ErrorCode.EMBEDDING_JOB_ATTEMPT_INVALID));
+
+ if (!Objects.equals(attempt.getId(), attemptId)
+ || attempt.getStatus() != AttemptStatus.STARTED
+ || attempt.getWorkerNode() == null
+ || !Objects.equals(attempt.getWorkerNode().getId(), workerId)) {
+ throw new DocGridException(ErrorCode.EMBEDDING_JOB_ATTEMPT_INVALID);
+ }
+ }
+
+ private void validateChunks(
+ DocumentVersion documentVersion,
+ List chunks
+ ) {
+ if (chunks == null || chunks.isEmpty()) {
+ throw new DocGridException(ErrorCode.DOCUMENT_CHUNKS_INCONSISTENT);
+ }
+
+ for (int index = 0; index < chunks.size(); index++) {
+ DocumentChunk chunk = chunks.get(index);
+ if (chunk == null
+ || chunk.getId() == null
+ || chunk.getDocumentVersion() == null
+ || !Objects.equals(chunk.getDocumentVersion().getId(), documentVersion.getId())
+ || chunk.getChunkIndex() != index
+ || !StringUtils.hasText(chunk.getChunkText())
+ || !StringUtils.hasText(chunk.getContentHash())
+ || !chunk.getContentHash().matches("[0-9a-f]{64}")) {
+ throw new DocGridException(ErrorCode.DOCUMENT_CHUNKS_INCONSISTENT);
+ }
+ }
+ }
+
+ private void validatePreparedTarget(
+ EmbeddingWork preparedWork,
+ DocumentVersion documentVersion,
+ EmbeddingModel embeddingModel
+ ) {
+ if (preparedWork == null
+ || !Objects.equals(preparedWork.documentVersionId(), documentVersion.getId())
+ || !Objects.equals(preparedWork.embeddingModelId(), embeddingModel.getId())
+ || preparedWork.dimension() != embeddingModel.getDimension()) {
+ throw new DocGridException(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT);
+ }
+ }
+
+ private void validatePreparedChunks(
+ EmbeddingWork preparedWork,
+ List chunks
+ ) {
+ if (preparedWork.chunks().size() != chunks.size()) {
+ throw new DocGridException(ErrorCode.DOCUMENT_CHUNKS_INCONSISTENT);
+ }
+
+ for (int index = 0; index < chunks.size(); index++) {
+ ChunkSnapshot snapshot = preparedWork.chunks().get(index);
+ DocumentChunk chunk = chunks.get(index);
+ if (snapshot == null
+ || !Objects.equals(snapshot.chunkId(), chunk.getId())
+ || snapshot.chunkIndex() != chunk.getChunkIndex()
+ || !Objects.equals(snapshot.chunkText(), chunk.getChunkText())
+ || !Objects.equals(snapshot.contentHash(), chunk.getContentHash())) {
+ throw new DocGridException(ErrorCode.DOCUMENT_CHUNKS_INCONSISTENT);
+ }
+ }
+ }
+
+ private EmbeddingState resolveState(
+ DocumentVersion documentVersion,
+ EmbeddingModel embeddingModel,
+ int chunkCount
+ ) {
+ long embeddingCount = embeddingRepository.countByDocumentVersionIdAndEmbeddingModelId(
+ documentVersion.getId(),
+ embeddingModel.getId()
+ );
+
+ if (documentVersion.getStatus() == DocumentVersionStatus.CHUNKED) {
+ if (embeddingCount != 0) {
+ throw new DocGridException(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT);
+ }
+ return EmbeddingState.WORK;
+ }
+ if (documentVersion.getStatus() == DocumentVersionStatus.EMBEDDING) {
+ if (embeddingCount == 0) {
+ return EmbeddingState.WORK;
+ }
+ if (embeddingCount == chunkCount) {
+ return EmbeddingState.REPLAY;
+ }
+ throw new DocGridException(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT);
+ }
+ if (embeddingCount != 0) {
+ throw new DocGridException(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT);
+ }
+ throw new DocGridException(ErrorCode.DOCUMENT_VERSION_EMBEDDING_NOT_ALLOWED);
+ }
+
+ private List toEntities(
+ DocumentVersion documentVersion,
+ EmbeddingModel embeddingModel,
+ List chunks,
+ List drafts
+ ) {
+ if (documentVersion.getDocument() == null
+ || documentVersion.getDocument().getId() == null
+ || drafts == null
+ || drafts.size() != chunks.size()) {
+ throw new DocGridException(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT);
+ }
+
+ List embeddings = new ArrayList<>(drafts.size());
+ for (int index = 0; index < drafts.size(); index++) {
+ DocumentChunk chunk = chunks.get(index);
+ DocumentEmbeddingDraft draft = drafts.get(index);
+ float[] vector = validateDraft(draft, chunk, embeddingModel.getDimension());
+ embeddings.add(Embedding.builder()
+ .chunk(chunk)
+ .document(documentVersion.getDocument())
+ .documentVersion(documentVersion)
+ .embeddingModel(embeddingModel)
+ .vector(vector)
+ .dimension(embeddingModel.getDimension())
+ .vectorHash(draft.vectorHash())
+ .status(EmbeddingStatus.ACTIVE)
+ .build());
+ }
+ return embeddings;
+ }
+
+ private float[] validateDraft(
+ DocumentEmbeddingDraft draft,
+ DocumentChunk chunk,
+ int expectedDimension
+ ) {
+ if (draft == null
+ || !Objects.equals(draft.chunkId(), chunk.getId())
+ || draft.chunkIndex() != chunk.getChunkIndex()
+ || !Objects.equals(draft.contentHash(), chunk.getContentHash())
+ || !StringUtils.hasText(draft.vectorHash())
+ || !draft.vectorHash().matches("[0-9a-f]{64}")) {
+ throw new DocGridException(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT);
+ }
+
+ float[] vector = draft.vector();
+ EmbeddingVectorSupport.validate(vector, expectedDimension);
+ if (!draft.vectorHash().equals(EmbeddingVectorSupport.calculateHash(vector))) {
+ throw new DocGridException(ErrorCode.EMBEDDING_VECTOR_INVALID);
+ }
+ return vector;
+ }
+
+ private CompletionResult result(
+ Long jobId,
+ Long attemptId,
+ DocumentVersion documentVersion,
+ EmbeddingModel embeddingModel,
+ int chunkCount,
+ boolean created
+ ) {
+ return new CompletionResult(
+ jobId,
+ attemptId,
+ documentVersion.getId(),
+ embeddingModel.getId(),
+ chunkCount,
+ chunkCount,
+ documentVersion.getStatus(),
+ created
+ );
+ }
+
+ private enum EmbeddingState {
+ WORK,
+ REPLAY
+ }
+
+ /**
+ * 외부 호출에 필요한 Version·Model 식별자와 정렬된 Chunk 값의 불변 Snapshot.
+ */
+ public record EmbeddingWork(
+ Long documentVersionId,
+ Long embeddingModelId,
+ int dimension,
+ List chunks
+ ) {
+
+ public EmbeddingWork {
+ chunks = List.copyOf(chunks);
+ }
+ }
+
+ /**
+ * 외부 호출에 전달하는 단일 Chunk의 식별자, 순서, Text와 내용 Hash Snapshot.
+ */
+ public record ChunkSnapshot(
+ Long chunkId,
+ int chunkIndex,
+ String chunkText,
+ String contentHash
+ ) {
+ }
+
+ /**
+ * 준비 Transaction이 전달하는 외부 작업 Snapshot 또는 기존 완료 결과 중 하나를 표현한다.
+ */
+ public record PreparationResult(
+ EmbeddingWork work,
+ CompletionResult replayResult
+ ) {
+
+ static PreparationResult work(EmbeddingWork work) {
+ return new PreparationResult(work, null);
+ }
+
+ static PreparationResult replay(CompletionResult replayResult) {
+ return new PreparationResult(null, replayResult);
+ }
+
+ public boolean isReplay() {
+ return replayResult != null;
+ }
+ }
+
+ /**
+ * 저장 생성 여부와 API 응답에 필요한 문서 Embedding Set 요약을 전달한다.
+ */
+ public record CompletionResult(
+ Long jobId,
+ Long attemptId,
+ Long documentVersionId,
+ Long embeddingModelId,
+ int chunkCount,
+ int embeddingCount,
+ DocumentVersionStatus documentVersionStatus,
+ boolean created
+ ) {
+ }
+}
diff --git a/src/main/java/com/opensource/docgrid/domain/embedding/service/query/QueryEmbeddingService.java b/src/main/java/com/opensource/docgrid/domain/embedding/service/query/QueryEmbeddingService.java
index 62c74c1..ea5453f 100644
--- a/src/main/java/com/opensource/docgrid/domain/embedding/service/query/QueryEmbeddingService.java
+++ b/src/main/java/com/opensource/docgrid/domain/embedding/service/query/QueryEmbeddingService.java
@@ -1,13 +1,9 @@
package com.opensource.docgrid.domain.embedding.service.query;
-import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.stereotype.Service;
-import org.springframework.web.client.RestClientException;
-import org.springframework.web.client.RestClient;
+import com.opensource.docgrid.domain.embedding.client.EmbeddingClient;
import com.opensource.docgrid.domain.embedding.dto.EmbedResult;
-import com.opensource.docgrid.domain.embedding.dto.request.EmbedRequest;
-import com.opensource.docgrid.domain.embedding.dto.response.EmbedServerResponse;
import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel;
import com.opensource.docgrid.global.exception.DocGridException;
import com.opensource.docgrid.global.exception.ErrorCode;
@@ -19,39 +15,26 @@
public class QueryEmbeddingService {
private final EmbeddingModelQueryService embeddingModelQueryService;
- private final RestClient restClient;
+ private final EmbeddingClient embeddingClient;
public QueryEmbeddingService(
EmbeddingModelQueryService embeddingModelQueryService,
- @Qualifier("embeddingRestClient") RestClient restClient
+ EmbeddingClient embeddingClient
) {
this.embeddingModelQueryService = embeddingModelQueryService;
- this.restClient = restClient;
+ this.embeddingClient = embeddingClient;
}
public EmbedResult embed(String text) {
EmbeddingModel activeModel = embeddingModelQueryService.getActiveModel();
+ float[] vector = embeddingClient.embed(text);
- EmbedServerResponse response;
- try {
- response = restClient.post()
- .uri("/embed")
- .body(new EmbedRequest(text))
- .retrieve()
- .body(EmbedServerResponse.class);
- } catch (RestClientException e) {
- log.error("임베딩 서버 호출 실패: {}", e.getMessage());
- throw new DocGridException(ErrorCode.EMBEDDING_SERVER_UNAVAILABLE);
- }
-
- if (response == null
- || response.vector() == null
- || response.vector().length != activeModel.getDimension()) {
- int actual = (response == null || response.vector() == null) ? -1 : response.vector().length;
+ if (vector == null || vector.length != activeModel.getDimension()) {
+ int actual = vector == null ? -1 : vector.length;
log.error("임베딩 차원 불일치: expected={}, actual={}", activeModel.getDimension(), actual);
throw new DocGridException(ErrorCode.EMBEDDING_DIMENSION_MISMATCH);
}
- return new EmbedResult(activeModel, response.vector());
+ return new EmbedResult(activeModel, vector);
}
}
diff --git a/src/main/java/com/opensource/docgrid/global/exception/ErrorCode.java b/src/main/java/com/opensource/docgrid/global/exception/ErrorCode.java
index d5137f3..f43ccbc 100644
--- a/src/main/java/com/opensource/docgrid/global/exception/ErrorCode.java
+++ b/src/main/java/com/opensource/docgrid/global/exception/ErrorCode.java
@@ -61,6 +61,9 @@ public enum ErrorCode {
DOCUMENT_VERSION_CHUNKING_NOT_ALLOWED(
HttpStatus.CONFLICT, "DOCUMENT-VERSION-005", "현재 문서 버전 상태에서는 Chunk를 생성할 수 없습니다."
),
+ DOCUMENT_VERSION_EMBEDDING_NOT_ALLOWED(
+ HttpStatus.CONFLICT, "DOCUMENT-VERSION-006", "현재 문서 버전 상태에서는 Embedding을 생성할 수 없습니다."
+ ),
INDEXING_STATUS_INCONSISTENT(
HttpStatus.INTERNAL_SERVER_ERROR, "DOCUMENT-STATUS-001", "문서 인덱싱 상태를 조회할 수 없습니다."
),
@@ -114,6 +117,16 @@ public enum ErrorCode {
"DOCUMENT-CHUNK-001",
"문서 버전과 Chunk 데이터가 일치하지 않습니다."
),
+ DOCUMENT_EMBEDDINGS_INCONSISTENT(
+ HttpStatus.INTERNAL_SERVER_ERROR,
+ "DOCUMENT-EMBEDDING-001",
+ "문서 버전과 Embedding 데이터가 일치하지 않습니다."
+ ),
+ EMBEDDING_VECTOR_INVALID(
+ HttpStatus.INTERNAL_SERVER_ERROR,
+ "DOCUMENT-EMBEDDING-002",
+ "생성된 Embedding Vector가 올바르지 않습니다."
+ ),
// PERMISSION
INVALID_TARGET_TYPE(HttpStatus.BAD_REQUEST, "PERMISSION-001", "target_type과 ID 필드 조합이 올바르지 않습니다."),
diff --git a/src/test/java/com/opensource/docgrid/domain/document/entity/DocumentVersionTest.java b/src/test/java/com/opensource/docgrid/domain/document/entity/DocumentVersionTest.java
index 2d435ce..67b77a4 100644
--- a/src/test/java/com/opensource/docgrid/domain/document/entity/DocumentVersionTest.java
+++ b/src/test/java/com/opensource/docgrid/domain/document/entity/DocumentVersionTest.java
@@ -13,7 +13,7 @@
import com.opensource.docgrid.domain.document.enums.DocumentVersionStatus;
/**
- * Document Version 파이프라인의 PARSING·CHUNKED 상태 전이 Guard와 기존 전이를 검증한다.
+ * Document Version 파이프라인의 PARSING·CHUNKED·EMBEDDING 상태 전이 Guard를 검증한다.
*
* Command Service를 우회한 잘못된 상태 변경은 즉시 실패하고 정상 순서만 허용되는지 확인한다.
*/
@@ -55,7 +55,28 @@ void markChunked_rejectsUnexpectedStatus(DocumentVersionStatus status) {
}
@Test
- @DisplayName("기존 EMBEDDING·INDEXED·FAILED 전이는 유지된다")
+ @DisplayName("CHUNKED에서 EMBEDDING으로 전이한다")
+ void markEmbedding_transitionsFromChunked() {
+ DocumentVersion version = version(DocumentVersionStatus.CHUNKED);
+
+ version.markEmbedding();
+
+ assertThat(version.getStatus()).isEqualTo(DocumentVersionStatus.EMBEDDING);
+ }
+
+ @ParameterizedTest
+ @EnumSource(value = DocumentVersionStatus.class, names = "CHUNKED", mode = EnumSource.Mode.EXCLUDE)
+ @DisplayName("CHUNKED가 아닌 상태에서는 EMBEDDING 전이를 거부한다")
+ void markEmbedding_rejectsUnexpectedStatus(DocumentVersionStatus status) {
+ DocumentVersion version = version(status);
+
+ assertThatThrownBy(version::markEmbedding)
+ .isInstanceOf(IllegalStateException.class);
+ assertThat(version.getStatus()).isEqualTo(status);
+ }
+
+ @Test
+ @DisplayName("EMBEDDING 이후 INDEXED·FAILED 전이는 유지된다")
void laterPipelineTransitions_arePreserved() {
DocumentVersion version = version(DocumentVersionStatus.CHUNKED);
LocalDateTime indexedAt = LocalDateTime.of(2026, 7, 29, 12, 0);
diff --git a/src/test/java/com/opensource/docgrid/domain/embedding/client/EmbeddingClientTest.java b/src/test/java/com/opensource/docgrid/domain/embedding/client/EmbeddingClientTest.java
new file mode 100644
index 0000000..d2422ad
--- /dev/null
+++ b/src/test/java/com/opensource/docgrid/domain/embedding/client/EmbeddingClientTest.java
@@ -0,0 +1,77 @@
+package com.opensource.docgrid.domain.embedding.client;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+import static org.mockito.BDDMockito.given;
+import static org.mockito.Mockito.doReturn;
+
+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.Answers;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+import org.mockito.junit.jupiter.MockitoSettings;
+import org.mockito.quality.Strictness;
+import org.springframework.web.client.ResourceAccessException;
+import org.springframework.web.client.RestClient;
+
+import com.opensource.docgrid.domain.embedding.dto.response.EmbedServerResponse;
+import com.opensource.docgrid.global.exception.DocGridException;
+import com.opensource.docgrid.global.exception.ErrorCode;
+
+/**
+ * EmbeddingClient의 HTTP 응답 전달과 외부 장애 변환 경계를 검증한다.
+ */
+@ExtendWith(MockitoExtension.class)
+@MockitoSettings(strictness = Strictness.LENIENT)
+@DisplayName("EmbeddingClient 단위 테스트")
+class EmbeddingClientTest {
+
+ @Mock private RestClient restClient;
+ @Mock(answer = Answers.RETURNS_SELF) private RestClient.RequestBodyUriSpec requestBodyUriSpec;
+ @Mock private RestClient.ResponseSpec responseSpec;
+
+ private EmbeddingClient embeddingClient;
+
+ @BeforeEach
+ void setUp() {
+ embeddingClient = new EmbeddingClient(restClient);
+ doReturn(requestBodyUriSpec).when(restClient).post();
+ doReturn(responseSpec).when(requestBodyUriSpec).retrieve();
+ }
+
+ @Test
+ @DisplayName("정상 응답: 외부 서버의 Vector를 그대로 반환한다")
+ void embed_success() {
+ float[] vector = new float[]{0.1f, 0.2f};
+ given(responseSpec.body(EmbedServerResponse.class))
+ .willReturn(new EmbedServerResponse(vector));
+
+ float[] result = embeddingClient.embed("검색어");
+
+ assertThat(result).containsExactly(vector);
+ }
+
+ @Test
+ @DisplayName("빈 응답: 외부 서버 응답이 null이면 null을 반환한다")
+ void embed_returnsNull_whenResponseIsNull() {
+ given(responseSpec.body(EmbedServerResponse.class)).willReturn(null);
+
+ float[] result = embeddingClient.embed("검색어");
+
+ assertThat(result).isNull();
+ }
+
+ @Test
+ @DisplayName("서버 장애: RestClientException을 서비스 사용 불가 오류로 변환한다")
+ void embed_throws_whenServerUnavailable() {
+ given(responseSpec.body(EmbedServerResponse.class))
+ .willThrow(new ResourceAccessException("Connection refused"));
+
+ assertThatThrownBy(() -> embeddingClient.embed("검색어"))
+ .isInstanceOf(DocGridException.class)
+ .hasFieldOrPropertyWithValue("errorCode", ErrorCode.EMBEDDING_SERVER_UNAVAILABLE);
+ }
+}
diff --git a/src/test/java/com/opensource/docgrid/domain/embedding/controller/IndexingJobAdminControllerTest.java b/src/test/java/com/opensource/docgrid/domain/embedding/controller/IndexingJobAdminControllerTest.java
index 55f6f43..5f87da2 100644
--- a/src/test/java/com/opensource/docgrid/domain/embedding/controller/IndexingJobAdminControllerTest.java
+++ b/src/test/java/com/opensource/docgrid/domain/embedding/controller/IndexingJobAdminControllerTest.java
@@ -32,8 +32,11 @@
import com.opensource.docgrid.domain.document.service.command.DocumentChunkTransactionService.ChunkResult;
import com.opensource.docgrid.domain.embedding.dto.response.ClaimedEmbeddingJobResponse;
import com.opensource.docgrid.domain.embedding.dto.response.DocumentChunksResponse;
+import com.opensource.docgrid.domain.embedding.dto.response.DocumentEmbeddingsResponse;
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.DocumentEmbeddingService.EmbeddingResult;
import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobAttemptService;
import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobAttemptService.StartResult;
import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobClaimService;
@@ -43,7 +46,7 @@
import com.opensource.docgrid.global.exception.ErrorCode;
/**
- * 관리자용 Embedding Job Claim, Attempt 시작과 Document Chunk 생성 API 계약을 검증하는 Web MVC 테스트.
+ * 관리자용 Job Claim, Attempt 시작과 Document Chunk·Embedding 생성 API 계약을 검증한다.
*
*
각 API의 최초 생성·멱등 재생·Validation·비즈니스 오류 및 ADMIN Security 동작을
* 실제 Service 실행 없이 Controller 경계에서 확인한다.
@@ -56,6 +59,7 @@ class IndexingJobAdminControllerTest {
private static final String CLAIM_URL = "/admin/indexing-jobs/claim";
private static final String ATTEMPT_URL = "/admin/indexing-jobs/10/attempts";
private static final String CHUNKS_URL = "/admin/indexing-jobs/10/attempts/100/chunks";
+ private static final String EMBEDDINGS_URL = "/admin/indexing-jobs/10/attempts/100/embeddings";
private static final Long JOB_ID = 10L;
private static final Long ATTEMPT_ID = 100L;
private static final Long WORKER_ID = 1L;
@@ -72,6 +76,7 @@ class IndexingJobAdminControllerTest {
@MockitoBean private EmbeddingJobClaimService embeddingJobClaimService;
@MockitoBean private EmbeddingJobAttemptService embeddingJobAttemptService;
@MockitoBean private DocumentParsingService documentParsingService;
+ @MockitoBean private DocumentEmbeddingService documentEmbeddingService;
@MockitoBean private JpaMetamodelMappingContext jpaMetamodelMappingContext;
@MockitoBean private JwtProvider jwtProvider;
@MockitoBean private CorsConfigurationSource corsConfigurationSource;
@@ -328,6 +333,95 @@ void createChunks_returnsForbidden_withoutAdminRole() throws Exception {
.andExpect(status().isForbidden());
}
+ @Test
+ @DisplayName("ADMIN 사용자의 최초 Embedding 저장은 201을 반환한다")
+ void createEmbeddings_returnsCreated_when_embeddingsAreCreated() throws Exception {
+ DocumentEmbeddingsResponse response = createEmbeddingsResponse();
+ given(documentEmbeddingService.createEmbeddings(eq(JOB_ID), eq(ATTEMPT_ID), any()))
+ .willReturn(new EmbeddingResult(response, true));
+
+ mockMvc.perform(post(EMBEDDINGS_URL)
+ .contentType("application/json")
+ .content(VALID_ATTEMPT_BODY)
+ .with(user("admin").roles("ADMIN")))
+ .andExpect(status().isCreated())
+ .andExpect(jsonPath("$.success").value(true))
+ .andExpect(jsonPath("$.data.jobId").value(JOB_ID))
+ .andExpect(jsonPath("$.data.attemptId").value(ATTEMPT_ID))
+ .andExpect(jsonPath("$.data.documentVersionId").value(5))
+ .andExpect(jsonPath("$.data.embeddingModelId").value(7))
+ .andExpect(jsonPath("$.data.chunkCount").value(3))
+ .andExpect(jsonPath("$.data.embeddingCount").value(3))
+ .andExpect(jsonPath("$.data.versionStatus").value("EMBEDDING"))
+ .andExpect(jsonPath("$.data.claimToken").doesNotExist())
+ .andExpect(jsonPath("$.data.vector").doesNotExist());
+ }
+
+ @Test
+ @DisplayName("ADMIN 사용자의 완료된 Embedding 재호출은 200을 반환한다")
+ void createEmbeddings_returnsOk_when_embeddingsAreReplayed() throws Exception {
+ DocumentEmbeddingsResponse response = createEmbeddingsResponse();
+ given(documentEmbeddingService.createEmbeddings(eq(JOB_ID), eq(ATTEMPT_ID), any()))
+ .willReturn(new EmbeddingResult(response, false));
+
+ mockMvc.perform(post(EMBEDDINGS_URL)
+ .contentType("application/json")
+ .content(VALID_ATTEMPT_BODY)
+ .with(user("admin").roles("ADMIN")))
+ .andExpect(status().isOk())
+ .andExpect(jsonPath("$.data.embeddingCount").value(3));
+ }
+
+ @ParameterizedTest(name = "{0}")
+ @MethodSource("invalidEmbeddingRequests")
+ @DisplayName("Embedding 생성 입력 형식이 올바르지 않으면 400을 반환한다")
+ void createEmbeddings_returnsBadRequest_when_requestIsInvalid(
+ String description,
+ String url,
+ String body
+ ) throws Exception {
+ mockMvc.perform(post(url)
+ .contentType("application/json")
+ .content(body)
+ .with(user("admin").roles("ADMIN")))
+ .andExpect(status().isBadRequest())
+ .andExpect(jsonPath("$.code").value("COMMON-002"));
+ }
+
+ @ParameterizedTest(name = "{0}")
+ @MethodSource("embeddingBusinessErrors")
+ @DisplayName("Embedding 생성 비즈니스 오류를 정의된 HTTP 상태와 코드로 반환한다")
+ void createEmbeddings_returnsDefinedError(
+ ErrorCode errorCode,
+ int expectedStatus,
+ String expectedCode
+ ) throws Exception {
+ given(documentEmbeddingService.createEmbeddings(eq(JOB_ID), eq(ATTEMPT_ID), any()))
+ .willThrow(new DocGridException(errorCode));
+
+ mockMvc.perform(post(EMBEDDINGS_URL)
+ .contentType("application/json")
+ .content(VALID_ATTEMPT_BODY)
+ .with(user("admin").roles("ADMIN")))
+ .andExpect(status().is(expectedStatus))
+ .andExpect(jsonPath("$.code").value(expectedCode));
+ }
+
+ @Test
+ @DisplayName("일반 사용자와 미인증 사용자는 Embedding을 생성할 수 없다")
+ void createEmbeddings_returnsForbidden_withoutAdminRole() throws Exception {
+ mockMvc.perform(post(EMBEDDINGS_URL)
+ .contentType("application/json")
+ .content(VALID_ATTEMPT_BODY)
+ .with(user("user").roles("USER")))
+ .andExpect(status().isForbidden());
+
+ mockMvc.perform(post(EMBEDDINGS_URL)
+ .contentType("application/json")
+ .content(VALID_ATTEMPT_BODY))
+ .andExpect(status().isForbidden());
+ }
+
private static Stream invalidAttemptRequests() {
return Stream.of(
Arguments.of("Job ID가 양수가 아님", "/admin/indexing-jobs/0/attempts", VALID_ATTEMPT_BODY),
@@ -389,6 +483,39 @@ private static Stream chunkBusinessErrors() {
);
}
+ private static Stream invalidEmbeddingRequests() {
+ return Stream.of(
+ Arguments.of(
+ "Job ID가 양수가 아님",
+ "/admin/indexing-jobs/0/attempts/100/embeddings",
+ VALID_ATTEMPT_BODY
+ ),
+ Arguments.of(
+ "Attempt ID가 양수가 아님",
+ "/admin/indexing-jobs/10/attempts/0/embeddings",
+ VALID_ATTEMPT_BODY
+ ),
+ Arguments.of("Worker ID가 양수가 아님", EMBEDDINGS_URL, """
+ {"workerId": 0, "claimToken": "%s"}
+ """.formatted(CLAIM_TOKEN)),
+ Arguments.of("Claim Token UUID 형식 오류", EMBEDDINGS_URL, """
+ {"workerId": 1, "claimToken": "not-a-uuid"}
+ """)
+ );
+ }
+
+ private static Stream embeddingBusinessErrors() {
+ return Stream.of(
+ Arguments.of(ErrorCode.EMBEDDING_JOB_NOT_FOUND, 404, "EMBEDDING-JOB-001"),
+ Arguments.of(ErrorCode.EMBEDDING_JOB_ATTEMPT_INVALID, 409, "EMBEDDING-JOB-006"),
+ Arguments.of(ErrorCode.DOCUMENT_VERSION_EMBEDDING_NOT_ALLOWED, 409, "DOCUMENT-VERSION-006"),
+ Arguments.of(ErrorCode.DOCUMENT_CHUNKS_INCONSISTENT, 500, "DOCUMENT-CHUNK-001"),
+ Arguments.of(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT, 500, "DOCUMENT-EMBEDDING-001"),
+ Arguments.of(ErrorCode.EMBEDDING_VECTOR_INVALID, 500, "DOCUMENT-EMBEDDING-002"),
+ Arguments.of(ErrorCode.EMBEDDING_SERVER_UNAVAILABLE, 503, "SEARCH-001")
+ );
+ }
+
private StartedEmbeddingJobAttemptResponse createAttemptResponse() {
return new StartedEmbeddingJobAttemptResponse(
100L,
@@ -409,4 +536,16 @@ private DocumentChunksResponse createChunksResponse() {
DocumentVersionStatus.CHUNKED
);
}
+
+ private DocumentEmbeddingsResponse createEmbeddingsResponse() {
+ return new DocumentEmbeddingsResponse(
+ JOB_ID,
+ ATTEMPT_ID,
+ 5L,
+ 7L,
+ 3,
+ 3,
+ DocumentVersionStatus.EMBEDDING
+ );
+ }
}
diff --git a/src/test/java/com/opensource/docgrid/domain/embedding/entity/EmbeddingTest.java b/src/test/java/com/opensource/docgrid/domain/embedding/entity/EmbeddingTest.java
new file mode 100644
index 0000000..bdb39ed
--- /dev/null
+++ b/src/test/java/com/opensource/docgrid/domain/embedding/entity/EmbeddingTest.java
@@ -0,0 +1,54 @@
+package com.opensource.docgrid.domain.embedding.entity;
+
+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.EmbeddingStatus;
+
+/**
+ * Embedding Entity가 Vector 배열과 기본 활성 상태를 불변 값으로 보존하는지 검증한다.
+ */
+@DisplayName("Embedding 테스트")
+class EmbeddingTest {
+
+ @Test
+ @DisplayName("생성에 사용한 Vector 배열을 변경해도 Entity 값은 유지된다")
+ void constructor_copiesVector() {
+ float[] vector = {0.1f, 0.2f};
+
+ Embedding embedding = Embedding.builder()
+ .vector(vector)
+ .dimension(vector.length)
+ .build();
+ vector[0] = 9.9f;
+
+ assertThat(embedding.getVector()).containsExactly(0.1f, 0.2f);
+ }
+
+ @Test
+ @DisplayName("조회한 Vector 배열을 변경해도 Entity 값은 유지된다")
+ void getter_returnsVectorCopy() {
+ Embedding embedding = Embedding.builder()
+ .vector(new float[]{0.1f, 0.2f})
+ .dimension(2)
+ .build();
+
+ float[] exposedVector = embedding.getVector();
+ exposedVector[0] = 9.9f;
+
+ assertThat(embedding.getVector()).containsExactly(0.1f, 0.2f);
+ }
+
+ @Test
+ @DisplayName("상태를 지정하지 않으면 ACTIVE로 생성된다")
+ void constructor_defaultsToActiveStatus() {
+ Embedding embedding = Embedding.builder()
+ .vector(new float[]{0.1f})
+ .dimension(1)
+ .build();
+
+ assertThat(embedding.getStatus()).isEqualTo(EmbeddingStatus.ACTIVE);
+ }
+}
diff --git a/src/test/java/com/opensource/docgrid/domain/embedding/integration/DocumentEmbeddingIntegrationTest.java b/src/test/java/com/opensource/docgrid/domain/embedding/integration/DocumentEmbeddingIntegrationTest.java
new file mode 100644
index 0000000..4897224
--- /dev/null
+++ b/src/test/java/com/opensource/docgrid/domain/embedding/integration/DocumentEmbeddingIntegrationTest.java
@@ -0,0 +1,403 @@
+package com.opensource.docgrid.domain.embedding.integration;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+
+import java.util.List;
+import java.util.Map;
+import java.util.UUID;
+import java.util.concurrent.CyclicBarrier;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+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.boot.test.context.TestConfiguration;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Import;
+import org.springframework.context.annotation.Primary;
+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 org.springframework.web.client.RestClient;
+
+import com.opensource.docgrid.domain.document.enums.DocumentVersionStatus;
+import com.opensource.docgrid.domain.embedding.client.EmbeddingClient;
+import com.opensource.docgrid.domain.embedding.dto.request.CreateDocumentEmbeddingsRequest;
+import com.opensource.docgrid.domain.embedding.service.DocumentEmbeddingService;
+import com.opensource.docgrid.domain.embedding.service.DocumentEmbeddingService.EmbeddingResult;
+import com.opensource.docgrid.global.exception.DocGridException;
+import com.opensource.docgrid.global.exception.ErrorCode;
+
+/**
+ * 실제 OpenSQL에서 Chunk Embedding의 Vector 저장, 실패 원자성과 동시 Transaction 수렴을 검증한다.
+ *
+ * 격리 Schema와 결정적인 외부 Client Stub을 사용해 Version·Model별 전체 Embedding Set만 저장되고,
+ * 같은 실행의 순차·동시 재호출이 하나의 결과로 수렴하는지 확인한다.
+ */
+@Tag("integration")
+@ActiveProfiles("test")
+@SpringBootTest
+@Import(DocumentEmbeddingIntegrationTest.EmbeddingClientTestConfig.class)
+@DirtiesContext(classMode = DirtiesContext.ClassMode.AFTER_CLASS)
+@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+@DisplayName("Document Embedding OpenSQL 통합 테스트")
+class DocumentEmbeddingIntegrationTest {
+
+ private static final String TEST_SCHEMA = "docgrid_document_embedding_test";
+ private static final int VECTOR_DIMENSION = 1024;
+ private static final int CONCURRENT_REQUESTS = 2;
+ private static final long TIMEOUT_SECONDS = 10;
+ private static final String CLAIM_TOKEN = "34c19d16-6ae1-4f6a-a35d-0123456789ab";
+ private static final String CONTENT_HASH =
+ "26e4a23eec4241e034f1b4631f0222f1895847637c35e77687d5945f75edb42c";
+
+ @Autowired private JdbcTemplate jdbcTemplate;
+ @Autowired private DocumentEmbeddingService documentEmbeddingService;
+ @Autowired private DeterministicEmbeddingClient embeddingClient;
+
+ @DynamicPropertySource
+ static void configureEmbedding(DynamicPropertyRegistry registry) {
+ registry.add("TEST_DB_SCHEMA", () -> TEST_SCHEMA);
+ registry.add("jwt.secret", () -> "docgrid-embedding-integration-test-secret-key-2026");
+ }
+
+ @BeforeEach
+ void resetState() {
+ embeddingClient.reset();
+ jdbcTemplate.execute("""
+ TRUNCATE TABLE
+ embeddings,
+ indexing_events,
+ document_chunks,
+ 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("Chunk 전체를 vector(1024)와 EMBEDDING 상태로 저장하고 재호출을 재생한다")
+ void createEmbeddings_persistsVectorSetAndReplays() {
+ ExecutionContext context = insertExecution();
+ CreateDocumentEmbeddingsRequest request =
+ new CreateDocumentEmbeddingsRequest(context.workerId(), CLAIM_TOKEN);
+
+ EmbeddingResult created =
+ documentEmbeddingService.createEmbeddings(context.jobId(), context.attemptId(), request);
+ EmbeddingResult replayed =
+ documentEmbeddingService.createEmbeddings(context.jobId(), context.attemptId(), request);
+
+ assertThat(created.created()).isTrue();
+ assertThat(replayed.created()).isFalse();
+ assertThat(replayed.response()).isEqualTo(created.response());
+ assertThat(created.response().embeddingCount()).isEqualTo(3);
+ assertThat(created.response().versionStatus()).isEqualTo(DocumentVersionStatus.EMBEDDING);
+ assertThat(embeddingClient.callCount()).isEqualTo(3);
+
+ assertThat(jdbcTemplate.queryForList("""
+ SELECT document_id, document_version_id, embedding_model_id, dimension, status,
+ vector_dims(vector) AS vector_dimension
+ FROM embeddings
+ WHERE document_version_id = ?
+ ORDER BY chunk_id
+ """, context.versionId()))
+ .containsExactly(
+ embeddingRow(context),
+ embeddingRow(context),
+ embeddingRow(context)
+ );
+ assertThat(jdbcTemplate.queryForList("""
+ SELECT vector_hash
+ FROM embeddings
+ WHERE document_version_id = ?
+ """, String.class, context.versionId()))
+ .allMatch(hash -> hash.matches("[0-9a-f]{64}"));
+ assertThat(queryString("SELECT status FROM document_versions WHERE id = ?", context.versionId()))
+ .isEqualTo("EMBEDDING");
+ assertThat(queryString("SELECT status FROM embedding_jobs WHERE id = ?", context.jobId()))
+ .isEqualTo("PROCESSING");
+ assertThat(queryString("SELECT status FROM embedding_job_attempts WHERE id = ?", context.attemptId()))
+ .isEqualTo("STARTED");
+ assertThat(eventCount(context.jobId(), "EMBEDDING_STARTED")).isOne();
+ }
+
+ @Test
+ @DisplayName("두 번째 Chunk 외부 호출이 실패하면 Embedding 행을 하나도 저장하지 않는다")
+ void createEmbeddings_externalFailureLeavesNoPartialRows() {
+ ExecutionContext context = insertExecution();
+ CreateDocumentEmbeddingsRequest request =
+ new CreateDocumentEmbeddingsRequest(context.workerId(), CLAIM_TOKEN);
+ embeddingClient.failOnText("두 번째");
+
+ assertThatThrownBy(() ->
+ documentEmbeddingService.createEmbeddings(context.jobId(), context.attemptId(), request))
+ .isInstanceOfSatisfying(DocGridException.class,
+ exception -> assertThat(exception.getErrorCode())
+ .isEqualTo(ErrorCode.EMBEDDING_SERVER_UNAVAILABLE));
+
+ assertThat(embeddingCount(context.versionId())).isZero();
+ assertThat(embeddingClient.callCount()).isEqualTo(2);
+ assertThat(queryString("SELECT status FROM document_versions WHERE id = ?", context.versionId()))
+ .isEqualTo("EMBEDDING");
+ assertThat(eventCount(context.jobId(), "EMBEDDING_STARTED")).isOne();
+ }
+
+ @Test
+ @DisplayName("같은 실행의 두 요청은 생성과 재생 하나씩 및 단일 Embedding Set으로 수렴한다")
+ void createEmbeddings_concurrentRequestsConverge() throws Exception {
+ ExecutionContext context = insertExecution();
+ CreateDocumentEmbeddingsRequest request =
+ new CreateDocumentEmbeddingsRequest(context.workerId(), CLAIM_TOKEN);
+ embeddingClient.armBarrier("첫 번째", CONCURRENT_REQUESTS);
+ ExecutorService executor = Executors.newFixedThreadPool(CONCURRENT_REQUESTS);
+
+ List results;
+ try {
+ List> futures = List.of(
+ executor.submit(() -> documentEmbeddingService.createEmbeddings(
+ context.jobId(), context.attemptId(), request
+ )),
+ executor.submit(() -> documentEmbeddingService.createEmbeddings(
+ context.jobId(), context.attemptId(), request
+ ))
+ );
+ results = List.of(
+ futures.get(0).get(TIMEOUT_SECONDS, TimeUnit.SECONDS),
+ futures.get(1).get(TIMEOUT_SECONDS, TimeUnit.SECONDS)
+ );
+ } finally {
+ embeddingClient.disarmBarrier();
+ executor.shutdownNow();
+ assertThat(executor.awaitTermination(TIMEOUT_SECONDS, TimeUnit.SECONDS)).isTrue();
+ }
+
+ assertThat(results).filteredOn(EmbeddingResult::created).hasSize(1);
+ assertThat(results).filteredOn(result -> !result.created()).hasSize(1);
+ assertThat(results).extracting(result -> result.response().embeddingCount()).containsOnly(3);
+ assertThat(embeddingCount(context.versionId())).isEqualTo(3);
+ assertThat(embeddingClient.callCount()).isEqualTo(6);
+ assertThat(eventCount(context.jobId(), "EMBEDDING_STARTED")).isOne();
+ }
+
+ private ExecutionContext insertExecution() {
+ String suffix = UUID.randomUUID().toString();
+ Long userId = jdbcTemplate.queryForObject("""
+ INSERT INTO users (email, password_hash, name, status, created_at, updated_at)
+ VALUES (?, 'password-hash', 'Embedding Test User', 'ACTIVE',
+ CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
+ RETURNING id
+ """, Long.class, "embedding-" + suffix + "@example.com");
+ Long workerId = jdbcTemplate.queryForObject("""
+ INSERT INTO worker_nodes (
+ worker_name, instance_id, status, last_heartbeat_at, started_at, created_at, updated_at
+ )
+ VALUES ('embedding-worker', ?, 'ACTIVE', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP,
+ CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
+ RETURNING id
+ """, Long.class, suffix);
+ Long documentId = jdbcTemplate.queryForObject("""
+ INSERT INTO documents (
+ owner_user_id, title, document_type, source_type, status, visibility, created_at, updated_at
+ )
+ VALUES (?, 'Embedding Test Document', 'TXT', 'UPLOAD', 'INDEXING', 'PRIVATE',
+ CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
+ RETURNING id
+ """, Long.class, userId);
+ 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, 'Embedding Test Version', 'text/plain', 'CHUNKED', ?,
+ CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
+ RETURNING id
+ """, Long.class, documentId, userId);
+ insertChunks(versionId);
+ Long embeddingModelId = jdbcTemplate.queryForObject("""
+ SELECT id
+ FROM embedding_models
+ WHERE is_active = TRUE AND is_searchable = TRUE
+ LIMIT 1
+ """, Long.class);
+ Long jobId = jdbcTemplate.queryForObject("""
+ INSERT INTO embedding_jobs (
+ document_version_id, embedding_model_id, status, priority, retry_count, max_retry_count,
+ locked_by_worker_id, locked_at, lock_expires_at, claim_token, started_at, created_at, updated_at
+ )
+ VALUES (?, ?, 'PROCESSING', 0, 0, 3, ?, CURRENT_TIMESTAMP,
+ TIMESTAMP '2099-01-01 00:00:00', ?, CURRENT_TIMESTAMP,
+ CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
+ RETURNING id
+ """, Long.class, versionId, embeddingModelId, workerId, CLAIM_TOKEN);
+ Long attemptId = jdbcTemplate.queryForObject("""
+ INSERT INTO embedding_job_attempts (
+ embedding_job_id, worker_node_id, attempt_no, claim_token, status,
+ started_at, created_at, updated_at
+ )
+ VALUES (?, ?, 1, ?, 'STARTED', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
+ RETURNING id
+ """, Long.class, jobId, workerId, CLAIM_TOKEN);
+ return new ExecutionContext(workerId, jobId, attemptId, documentId, versionId, embeddingModelId);
+ }
+
+ private void insertChunks(Long versionId) {
+ jdbcTemplate.update("""
+ INSERT INTO document_chunks (
+ document_version_id, chunk_index, chunk_text, token_count, char_start, char_end,
+ content_hash, created_at, updated_at
+ )
+ VALUES
+ (?, 0, '첫 번째', 1, 0, 4, ?, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
+ (?, 1, '두 번째', 1, 4, 8, ?, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
+ (?, 2, '세 번째', 1, 8, 12, ?, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
+ """,
+ versionId, CONTENT_HASH,
+ versionId, CONTENT_HASH,
+ versionId, CONTENT_HASH
+ );
+ }
+
+ private Map embeddingRow(ExecutionContext context) {
+ return Map.of(
+ "document_id", context.documentId(),
+ "document_version_id", context.versionId(),
+ "embedding_model_id", context.embeddingModelId(),
+ "dimension", VECTOR_DIMENSION,
+ "status", "ACTIVE",
+ "vector_dimension", VECTOR_DIMENSION
+ );
+ }
+
+ private int embeddingCount(Long versionId) {
+ return jdbcTemplate.queryForObject(
+ "SELECT COUNT(*) FROM embeddings WHERE document_version_id = ?",
+ Integer.class,
+ versionId
+ );
+ }
+
+ private String queryString(String sql, Long id) {
+ return jdbcTemplate.queryForObject(sql, String.class, id);
+ }
+
+ private int eventCount(Long jobId, String eventType) {
+ return jdbcTemplate.queryForObject(
+ "SELECT COUNT(*) FROM indexing_events WHERE embedding_job_id = ? AND event_type = ?",
+ Integer.class,
+ jobId,
+ eventType
+ );
+ }
+
+ /**
+ * 통합 테스트 실행 Context의 Worker, Job, Attempt, Document, Version과 Model 식별자를 묶는다.
+ */
+ private record ExecutionContext(
+ Long workerId,
+ Long jobId,
+ Long attemptId,
+ Long documentId,
+ Long versionId,
+ Long embeddingModelId
+ ) {
+ }
+
+ /**
+ * 실제 DB Transaction 밖의 외부 호출을 결정적으로 제어하는 테스트 Client Bean을 제공한다.
+ */
+ @TestConfiguration
+ static class EmbeddingClientTestConfig {
+
+ @Bean
+ @Primary
+ DeterministicEmbeddingClient deterministicEmbeddingClient() {
+ return new DeterministicEmbeddingClient();
+ }
+ }
+
+ /**
+ * Text별 1024차원 Vector를 만들고 실패 지점과 두 요청의 동시 진입 Barrier를 제어한다.
+ */
+ static class DeterministicEmbeddingClient extends EmbeddingClient {
+
+ private final AtomicInteger calls = new AtomicInteger();
+ private volatile String failureText;
+ private volatile String barrierText;
+ private volatile CyclicBarrier barrier;
+
+ DeterministicEmbeddingClient() {
+ super(RestClient.builder().build());
+ }
+
+ @Override
+ public float[] embed(String text) {
+ calls.incrementAndGet();
+ if (text.equals(failureText)) {
+ throw new DocGridException(ErrorCode.EMBEDDING_SERVER_UNAVAILABLE);
+ }
+
+ CyclicBarrier currentBarrier = barrier;
+ if (currentBarrier != null && text.equals(barrierText)) {
+ try {
+ currentBarrier.await(TIMEOUT_SECONDS, TimeUnit.SECONDS);
+ } catch (Exception exception) {
+ throw new IllegalStateException("Embedding 통합 테스트 Barrier 대기에 실패했습니다.", exception);
+ }
+ }
+
+ float[] vector = new float[VECTOR_DIMENSION];
+ vector[0] = text.hashCode();
+ vector[1] = text.length();
+ return vector;
+ }
+
+ int callCount() {
+ return calls.get();
+ }
+
+ void failOnText(String text) {
+ failureText = text;
+ }
+
+ void armBarrier(String text, int parties) {
+ barrierText = text;
+ barrier = new CyclicBarrier(parties);
+ }
+
+ void disarmBarrier() {
+ CyclicBarrier currentBarrier = barrier;
+ barrier = null;
+ barrierText = null;
+ if (currentBarrier != null) {
+ currentBarrier.reset();
+ }
+ }
+
+ void reset() {
+ disarmBarrier();
+ failureText = null;
+ calls.set(0);
+ }
+ }
+}
diff --git a/src/test/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingGeneratorTest.java b/src/test/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingGeneratorTest.java
new file mode 100644
index 0000000..8d23173
--- /dev/null
+++ b/src/test/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingGeneratorTest.java
@@ -0,0 +1,128 @@
+package com.opensource.docgrid.domain.embedding.service;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.assertj.core.api.Assertions.assertThatThrownBy;
+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.util.List;
+
+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.embedding.client.EmbeddingClient;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.ChunkSnapshot;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.EmbeddingWork;
+import com.opensource.docgrid.global.exception.DocGridException;
+import com.opensource.docgrid.global.exception.ErrorCode;
+
+/**
+ * Chunk Vector의 단건 순차 생성, 검증 실패 중단, Hash와 Draft 불변성 계약을 검증한다.
+ */
+@ExtendWith(MockitoExtension.class)
+@DisplayName("DocumentEmbeddingGenerator 테스트")
+class DocumentEmbeddingGeneratorTest {
+
+ private static final String CONTENT_HASH =
+ "26e4a23eec4241e034f1b4631f0222f1895847637c35e77687d5945f75edb42c";
+ private static final String VECTOR_HASH =
+ "efde6f1eea3d119bf25d4a0e2ce188572e80986e14d85c3f2a9528fa66f8731f";
+
+ @Mock private EmbeddingClient embeddingClient;
+
+ @Test
+ @DisplayName("Chunk 순서대로 단건 호출하고 검증된 Vector와 SHA-256 Hash를 반환한다")
+ void generate_callsSequentiallyAndReturnsDrafts() {
+ DocumentEmbeddingGenerator generator = new DocumentEmbeddingGenerator(embeddingClient);
+ given(embeddingClient.embed("첫 번째")).willReturn(new float[]{1.0f, -2.0f});
+ given(embeddingClient.embed("두 번째")).willReturn(new float[]{0.5f, 0.25f});
+
+ List drafts = generator.generate(work(
+ chunk(20L, 0, "첫 번째"),
+ chunk(21L, 1, "두 번째")
+ ));
+
+ InOrder callOrder = inOrder(embeddingClient);
+ callOrder.verify(embeddingClient).embed("첫 번째");
+ callOrder.verify(embeddingClient).embed("두 번째");
+ assertThat(drafts)
+ .extracting(DocumentEmbeddingDraft::chunkId)
+ .containsExactly(20L, 21L);
+ assertThat(drafts.get(0).vector()).containsExactly(1.0f, -2.0f);
+ assertThat(drafts.get(0).vectorHash()).isEqualTo(VECTOR_HASH);
+ }
+
+ @Test
+ @DisplayName("첫 Vector 차원이 다르면 뒤 Chunk를 호출하지 않고 실패한다")
+ void generate_stopsWhenDimensionMismatches() {
+ DocumentEmbeddingGenerator generator = new DocumentEmbeddingGenerator(embeddingClient);
+ given(embeddingClient.embed("첫 번째")).willReturn(new float[]{1.0f});
+
+ assertThatThrownBy(() -> generator.generate(work(
+ chunk(20L, 0, "첫 번째"),
+ chunk(21L, 1, "두 번째")
+ )))
+ .isInstanceOfSatisfying(DocGridException.class,
+ exception -> assertThat(exception.getErrorCode())
+ .isEqualTo(ErrorCode.EMBEDDING_DIMENSION_MISMATCH));
+
+ then(embeddingClient).should(never()).embed("두 번째");
+ }
+
+ @Test
+ @DisplayName("NaN이 포함된 Vector를 저장 Draft로 만들지 않는다")
+ void generate_rejectsNonFiniteVector() {
+ DocumentEmbeddingGenerator generator = new DocumentEmbeddingGenerator(embeddingClient);
+ given(embeddingClient.embed("첫 번째")).willReturn(new float[]{Float.NaN, 1.0f});
+
+ assertThatThrownBy(() -> generator.generate(work(chunk(20L, 0, "첫 번째"))))
+ .isInstanceOfSatisfying(DocGridException.class,
+ exception -> assertThat(exception.getErrorCode())
+ .isEqualTo(ErrorCode.EMBEDDING_VECTOR_INVALID));
+ }
+
+ @Test
+ @DisplayName("빈 작업 Snapshot은 외부 서버 호출 전에 거부한다")
+ void generate_rejectsEmptyWork() {
+ DocumentEmbeddingGenerator generator = new DocumentEmbeddingGenerator(embeddingClient);
+ EmbeddingWork emptyWork = new EmbeddingWork(5L, 7L, 2, List.of());
+
+ assertThatThrownBy(() -> generator.generate(emptyWork))
+ .isInstanceOfSatisfying(DocGridException.class,
+ exception -> assertThat(exception.getErrorCode())
+ .isEqualTo(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT));
+
+ then(embeddingClient).shouldHaveNoInteractions();
+ }
+
+ @Test
+ @DisplayName("Draft에서 조회한 Vector를 변경해도 보관 값은 유지된다")
+ void draft_returnsVectorCopy() {
+ DocumentEmbeddingDraft draft = new DocumentEmbeddingDraft(
+ 20L,
+ 0,
+ CONTENT_HASH,
+ new float[]{1.0f, -2.0f},
+ VECTOR_HASH
+ );
+
+ float[] exposedVector = draft.vector();
+ exposedVector[0] = 9.9f;
+
+ assertThat(draft.vector()).containsExactly(1.0f, -2.0f);
+ }
+
+ private EmbeddingWork work(ChunkSnapshot... chunks) {
+ return new EmbeddingWork(5L, 7L, 2, List.of(chunks));
+ }
+
+ private ChunkSnapshot chunk(Long chunkId, int chunkIndex, String text) {
+ return new ChunkSnapshot(chunkId, chunkIndex, text, CONTENT_HASH);
+ }
+}
diff --git a/src/test/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingServiceTest.java b/src/test/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingServiceTest.java
new file mode 100644
index 0000000..73b675f
--- /dev/null
+++ b/src/test/java/com/opensource/docgrid/domain/embedding/service/DocumentEmbeddingServiceTest.java
@@ -0,0 +1,154 @@
+package com.opensource.docgrid.domain.embedding.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.util.List;
+
+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.embedding.dto.request.CreateDocumentEmbeddingsRequest;
+import com.opensource.docgrid.domain.embedding.service.DocumentEmbeddingService.EmbeddingResult;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.ChunkSnapshot;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.CompletionResult;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.EmbeddingWork;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.PreparationResult;
+import com.opensource.docgrid.global.exception.DocGridException;
+import com.opensource.docgrid.global.exception.ErrorCode;
+
+/**
+ * 준비·외부 Vector 생성·완료 Transaction의 실행 순서와 재생·실패 중단 계약을 검증한다.
+ */
+@ExtendWith(MockitoExtension.class)
+@DisplayName("DocumentEmbeddingService 테스트")
+class DocumentEmbeddingServiceTest {
+
+ private static final Long JOB_ID = 10L;
+ private static final Long ATTEMPT_ID = 100L;
+ private static final Long VERSION_ID = 5L;
+ private static final Long MODEL_ID = 7L;
+ private static final Long WORKER_ID = 1L;
+ private static final String CLAIM_TOKEN = "34c19d16-6ae1-4f6a-a35d-0123456789ab";
+ private static final String CONTENT_HASH =
+ "26e4a23eec4241e034f1b4631f0222f1895847637c35e77687d5945f75edb42c";
+
+ @Mock private DocumentEmbeddingTransactionService transactionService;
+ @Mock private DocumentEmbeddingGenerator embeddingGenerator;
+
+ private DocumentEmbeddingService service;
+ private CreateDocumentEmbeddingsRequest request;
+ private EmbeddingWork work;
+
+ @BeforeEach
+ void setUp() {
+ service = new DocumentEmbeddingService(transactionService, embeddingGenerator);
+ request = new CreateDocumentEmbeddingsRequest(WORKER_ID, CLAIM_TOKEN);
+ work = new EmbeddingWork(
+ VERSION_ID,
+ MODEL_ID,
+ 2,
+ List.of(new ChunkSnapshot(20L, 0, "본문", CONTENT_HASH))
+ );
+ }
+
+ @Test
+ @DisplayName("준비, Vector 생성, 완료 순서로 실행하고 최초 저장 결과를 반환한다")
+ void createEmbeddings_executesExternalWorkBetweenTransactions() {
+ List drafts = List.of(draft());
+ CompletionResult completion = completion(true);
+ given(transactionService.prepare(JOB_ID, ATTEMPT_ID, WORKER_ID, CLAIM_TOKEN))
+ .willReturn(new PreparationResult(work, null));
+ given(embeddingGenerator.generate(work)).willReturn(drafts);
+ given(transactionService.complete(
+ JOB_ID,
+ ATTEMPT_ID,
+ WORKER_ID,
+ CLAIM_TOKEN,
+ work,
+ drafts
+ )).willReturn(completion);
+
+ EmbeddingResult result = service.createEmbeddings(JOB_ID, ATTEMPT_ID, request);
+
+ assertThat(result.created()).isTrue();
+ assertThat(result.response().documentVersionId()).isEqualTo(VERSION_ID);
+ assertThat(result.response().embeddingModelId()).isEqualTo(MODEL_ID);
+ assertThat(result.response().embeddingCount()).isOne();
+ assertThat(result.response().versionStatus()).isEqualTo(DocumentVersionStatus.EMBEDDING);
+
+ InOrder order = inOrder(transactionService, embeddingGenerator);
+ order.verify(transactionService).prepare(JOB_ID, ATTEMPT_ID, WORKER_ID, CLAIM_TOKEN);
+ order.verify(embeddingGenerator).generate(work);
+ order.verify(transactionService).complete(
+ JOB_ID, ATTEMPT_ID, WORKER_ID, CLAIM_TOKEN, work, drafts
+ );
+ }
+
+ @Test
+ @DisplayName("준비 단계가 완료 결과를 반환하면 외부 호출과 완료 Transaction 없이 재생한다")
+ void createEmbeddings_replaysWithoutExternalWork() {
+ CompletionResult replay = completion(false);
+ given(transactionService.prepare(JOB_ID, ATTEMPT_ID, WORKER_ID, CLAIM_TOKEN))
+ .willReturn(new PreparationResult(null, replay));
+
+ EmbeddingResult result = service.createEmbeddings(JOB_ID, ATTEMPT_ID, request);
+
+ assertThat(result.created()).isFalse();
+ assertThat(result.response().embeddingCount()).isOne();
+ then(embeddingGenerator).shouldHaveNoInteractions();
+ then(transactionService).should(never()).complete(any(), any(), any(), any(), any(), any());
+ }
+
+ @Test
+ @DisplayName("외부 Vector 생성이 실패하면 완료 Transaction을 호출하지 않는다")
+ void createEmbeddings_stopsWhenGeneratorFails() {
+ given(transactionService.prepare(JOB_ID, ATTEMPT_ID, WORKER_ID, CLAIM_TOKEN))
+ .willReturn(new PreparationResult(work, null));
+ given(embeddingGenerator.generate(work))
+ .willThrow(new DocGridException(ErrorCode.EMBEDDING_SERVER_UNAVAILABLE));
+
+ assertThatThrownBy(() -> service.createEmbeddings(JOB_ID, ATTEMPT_ID, request))
+ .isInstanceOfSatisfying(DocGridException.class,
+ exception -> assertThat(exception.getErrorCode())
+ .isEqualTo(ErrorCode.EMBEDDING_SERVER_UNAVAILABLE));
+
+ then(transactionService).should(never()).complete(any(), any(), any(), any(), any(), any());
+ }
+
+ private DocumentEmbeddingDraft draft() {
+ float[] vector = {0.1f, 0.2f};
+ return new DocumentEmbeddingDraft(
+ 20L,
+ 0,
+ CONTENT_HASH,
+ vector,
+ EmbeddingVectorSupport.calculateHash(vector)
+ );
+ }
+
+ private CompletionResult completion(boolean created) {
+ return new CompletionResult(
+ JOB_ID,
+ ATTEMPT_ID,
+ VERSION_ID,
+ MODEL_ID,
+ 1,
+ 1,
+ DocumentVersionStatus.EMBEDDING,
+ created
+ );
+ }
+}
diff --git a/src/test/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionServiceTest.java b/src/test/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionServiceTest.java
new file mode 100644
index 0000000..5ee3d0b
--- /dev/null
+++ b/src/test/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionServiceTest.java
@@ -0,0 +1,467 @@
+package com.opensource.docgrid.domain.embedding.service.command;
+
+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.never;
+
+import java.time.Clock;
+import java.time.Instant;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+import java.util.List;
+import java.util.Optional;
+
+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 org.springframework.test.util.ReflectionTestUtils;
+
+import com.opensource.docgrid.domain.document.entity.Document;
+import com.opensource.docgrid.domain.document.entity.DocumentChunk;
+import com.opensource.docgrid.domain.document.entity.DocumentVersion;
+import com.opensource.docgrid.domain.document.enums.DocumentSourceType;
+import com.opensource.docgrid.domain.document.enums.DocumentStatus;
+import com.opensource.docgrid.domain.document.enums.DocumentType;
+import com.opensource.docgrid.domain.document.enums.DocumentVersionStatus;
+import com.opensource.docgrid.domain.document.enums.VisibilityType;
+import com.opensource.docgrid.domain.document.repository.DocumentChunkRepository;
+import com.opensource.docgrid.domain.document.repository.DocumentVersionRepository;
+import com.opensource.docgrid.domain.embedding.entity.EmbeddingJob;
+import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel;
+import com.opensource.docgrid.domain.embedding.entity.Embedding;
+import com.opensource.docgrid.domain.embedding.enums.EmbeddingStatus;
+import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus;
+import com.opensource.docgrid.domain.embedding.fixture.EmbeddingModelFixture;
+import com.opensource.docgrid.domain.embedding.repository.EmbeddingJobRepository;
+import com.opensource.docgrid.domain.embedding.repository.EmbeddingRepository;
+import com.opensource.docgrid.domain.embedding.service.DocumentEmbeddingDraft;
+import com.opensource.docgrid.domain.embedding.service.EmbeddingVectorSupport;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.ChunkSnapshot;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.CompletionResult;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.EmbeddingWork;
+import com.opensource.docgrid.domain.embedding.service.command.DocumentEmbeddingTransactionService.PreparationResult;
+import com.opensource.docgrid.domain.worker.entity.EmbeddingJobAttempt;
+import com.opensource.docgrid.domain.worker.entity.IndexingEvent;
+import com.opensource.docgrid.domain.worker.entity.WorkerNode;
+import com.opensource.docgrid.domain.worker.enums.AttemptStatus;
+import com.opensource.docgrid.domain.worker.enums.IndexingEventType;
+import com.opensource.docgrid.domain.worker.enums.WorkerStatus;
+import com.opensource.docgrid.domain.worker.repository.EmbeddingJobAttemptRepository;
+import com.opensource.docgrid.domain.worker.repository.IndexingEventRepository;
+import com.opensource.docgrid.global.exception.DocGridException;
+import com.opensource.docgrid.global.exception.ErrorCode;
+
+/**
+ * Chunk Embedding 준비·완료 Transaction의 상태 전이, 재개, 완료 재생과 원자 저장 계약을 검증한다.
+ */
+@ExtendWith(MockitoExtension.class)
+@DisplayName("DocumentEmbeddingTransactionService 테스트")
+class DocumentEmbeddingTransactionServiceTest {
+
+ private static final Long JOB_ID = 10L;
+ private static final Long ATTEMPT_ID = 100L;
+ private static final Long VERSION_ID = 5L;
+ private static final Long MODEL_ID = 7L;
+ private static final Long WORKER_ID = 1L;
+ private static final String CLAIM_TOKEN = "34c19d16-6ae1-4f6a-a35d-0123456789ab";
+ private static final String CONTENT_HASH =
+ "26e4a23eec4241e034f1b4631f0222f1895847637c35e77687d5945f75edb42c";
+ private static final LocalDateTime NOW = LocalDateTime.of(2026, 7, 29, 12, 0);
+
+ @Mock private EmbeddingJobRepository embeddingJobRepository;
+ @Mock private EmbeddingJobAttemptRepository embeddingJobAttemptRepository;
+ @Mock private DocumentVersionRepository documentVersionRepository;
+ @Mock private DocumentChunkRepository documentChunkRepository;
+ @Mock private EmbeddingRepository embeddingRepository;
+ @Mock private IndexingEventRepository indexingEventRepository;
+ @Mock private EmbeddingJobOwnershipValidator ownershipValidator;
+
+ private DocumentEmbeddingTransactionService service;
+ private EmbeddingJob embeddingJob;
+ private DocumentVersion documentVersion;
+ private EmbeddingModel embeddingModel;
+ private EmbeddingJobAttempt attempt;
+
+ @BeforeEach
+ void setUp() {
+ Clock clock = Clock.fixed(
+ Instant.parse("2026-07-29T03:00:00Z"),
+ ZoneId.of("Asia/Seoul")
+ );
+ service = new DocumentEmbeddingTransactionService(
+ embeddingJobRepository,
+ embeddingJobAttemptRepository,
+ documentVersionRepository,
+ documentChunkRepository,
+ embeddingRepository,
+ indexingEventRepository,
+ ownershipValidator,
+ clock
+ );
+ prepareEntities(DocumentVersionStatus.CHUNKED);
+ }
+
+ @Test
+ @DisplayName("CHUNKED Version을 EMBEDDING으로 전환하고 고정 Model의 Chunk Snapshot을 만든다")
+ void prepare_marksEmbeddingAndReturnsSnapshot() {
+ givenValidContext(List.of(chunk(0, "첫 번째"), chunk(1, "두 번째")));
+
+ PreparationResult result = service.prepare(JOB_ID, ATTEMPT_ID, WORKER_ID, CLAIM_TOKEN);
+
+ assertThat(result.isReplay()).isFalse();
+ assertThat(result.work().documentVersionId()).isEqualTo(VERSION_ID);
+ assertThat(result.work().embeddingModelId()).isEqualTo(MODEL_ID);
+ assertThat(result.work().dimension()).isEqualTo(EmbeddingModelFixture.DIMENSION);
+ assertThat(result.work().chunks())
+ .extracting(chunk -> chunk.chunkIndex() + ":" + chunk.chunkText())
+ .containsExactly("0:첫 번째", "1:두 번째");
+ assertThat(documentVersion.getStatus()).isEqualTo(DocumentVersionStatus.EMBEDDING);
+ then(ownershipValidator).should().validate(embeddingJob, WORKER_ID, CLAIM_TOKEN, NOW);
+
+ ArgumentCaptor eventCaptor = ArgumentCaptor.forClass(IndexingEvent.class);
+ then(indexingEventRepository).should().save(eventCaptor.capture());
+ assertThat(eventCaptor.getValue().getEventType()).isEqualTo(IndexingEventType.EMBEDDING_STARTED);
+ }
+
+ @Test
+ @DisplayName("저장 결과가 없는 EMBEDDING Version은 시작 이벤트 없이 작업을 재개한다")
+ void prepare_resumesEmbeddingWithoutDuplicateEvent() {
+ prepareEntities(DocumentVersionStatus.EMBEDDING);
+ givenValidContext(List.of(chunk(0, "본문")));
+
+ PreparationResult result = service.prepare(JOB_ID, ATTEMPT_ID, WORKER_ID, CLAIM_TOKEN);
+
+ assertThat(result.isReplay()).isFalse();
+ assertThat(result.work().chunks()).hasSize(1);
+ then(indexingEventRepository).shouldHaveNoInteractions();
+ }
+
+ @Test
+ @DisplayName("모든 Chunk의 Embedding이 저장됐으면 기존 완료 결과를 재생한다")
+ void prepare_replaysCompletedEmbeddings() {
+ prepareEntities(DocumentVersionStatus.EMBEDDING);
+ givenValidContext(List.of(chunk(0, "첫 번째"), chunk(1, "두 번째")));
+ given(embeddingRepository.countByDocumentVersionIdAndEmbeddingModelId(VERSION_ID, MODEL_ID))
+ .willReturn(2L);
+
+ PreparationResult result = service.prepare(JOB_ID, ATTEMPT_ID, WORKER_ID, CLAIM_TOKEN);
+
+ assertThat(result.isReplay()).isTrue();
+ assertThat(result.work()).isNull();
+ assertThat(result.replayResult().created()).isFalse();
+ assertThat(result.replayResult().embeddingCount()).isEqualTo(2);
+ then(indexingEventRepository).shouldHaveNoInteractions();
+ }
+
+ @Test
+ @DisplayName("일부 Chunk의 Embedding만 저장된 상태는 내부 데이터 모순으로 거부한다")
+ void prepare_rejectsPartialEmbeddings() {
+ prepareEntities(DocumentVersionStatus.EMBEDDING);
+ givenValidContext(List.of(chunk(0, "첫 번째"), chunk(1, "두 번째")));
+ given(embeddingRepository.countByDocumentVersionIdAndEmbeddingModelId(VERSION_ID, MODEL_ID))
+ .willReturn(1L);
+
+ assertThatThrownBy(() -> service.prepare(JOB_ID, ATTEMPT_ID, WORKER_ID, CLAIM_TOKEN))
+ .isInstanceOfSatisfying(DocGridException.class,
+ exception -> assertThat(exception.getErrorCode())
+ .isEqualTo(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT));
+ }
+
+ @Test
+ @DisplayName("현재 Claim과 일치하지 않는 Attempt는 Version 조회 전에 거부한다")
+ void prepare_rejectsInvalidAttempt() {
+ given(embeddingJobRepository.findByIdForUpdate(JOB_ID)).willReturn(Optional.of(embeddingJob));
+ given(embeddingJobAttemptRepository.findByEmbeddingJobIdAndClaimToken(JOB_ID, CLAIM_TOKEN))
+ .willReturn(Optional.empty());
+
+ assertThatThrownBy(() -> service.prepare(JOB_ID, ATTEMPT_ID, WORKER_ID, CLAIM_TOKEN))
+ .isInstanceOfSatisfying(DocGridException.class,
+ exception -> assertThat(exception.getErrorCode())
+ .isEqualTo(ErrorCode.EMBEDDING_JOB_ATTEMPT_INVALID));
+
+ then(documentVersionRepository).shouldHaveNoInteractions();
+ then(documentChunkRepository).shouldHaveNoInteractions();
+ }
+
+ @Test
+ @DisplayName("Chunk Index가 연속되지 않으면 외부 호출 Snapshot을 만들지 않는다")
+ void prepare_rejectsNonSequentialChunks() {
+ givenValidContext(List.of(chunk(1, "본문")));
+
+ assertThatThrownBy(() -> service.prepare(JOB_ID, ATTEMPT_ID, WORKER_ID, CLAIM_TOKEN))
+ .isInstanceOfSatisfying(DocGridException.class,
+ exception -> assertThat(exception.getErrorCode())
+ .isEqualTo(ErrorCode.DOCUMENT_CHUNKS_INCONSISTENT));
+
+ then(embeddingRepository).shouldHaveNoInteractions();
+ }
+
+ @Test
+ @DisplayName("CHUNKED나 EMBEDDING이 아닌 Version은 Embedding 준비를 거부한다")
+ void prepare_rejectsUnexpectedVersionStatus() {
+ prepareEntities(DocumentVersionStatus.INDEXED);
+ givenValidContext(List.of(chunk(0, "본문")));
+
+ assertThatThrownBy(() -> service.prepare(JOB_ID, ATTEMPT_ID, WORKER_ID, CLAIM_TOKEN))
+ .isInstanceOfSatisfying(DocGridException.class,
+ exception -> assertThat(exception.getErrorCode())
+ .isEqualTo(ErrorCode.DOCUMENT_VERSION_EMBEDDING_NOT_ALLOWED));
+ }
+
+ @Test
+ @DisplayName("검증된 Draft 전체를 ACTIVE Embedding Set으로 저장하고 EMBEDDING 상태를 유지한다")
+ @SuppressWarnings("unchecked")
+ void complete_savesEmbeddingSet() {
+ prepareEntities(DocumentVersionStatus.EMBEDDING);
+ List chunks = List.of(chunk(0, "첫 번째"), chunk(1, "두 번째"));
+ givenValidContext(chunks);
+ EmbeddingWork work = work(chunks);
+ List drafts = List.of(
+ draft(chunks.get(0), 0.1f),
+ draft(chunks.get(1), 0.2f)
+ );
+
+ CompletionResult result = service.complete(
+ JOB_ID,
+ ATTEMPT_ID,
+ WORKER_ID,
+ CLAIM_TOKEN,
+ work,
+ drafts
+ );
+
+ assertThat(result.created()).isTrue();
+ assertThat(result.embeddingCount()).isEqualTo(2);
+ assertThat(result.documentVersionStatus()).isEqualTo(DocumentVersionStatus.EMBEDDING);
+ assertThat(documentVersion.getStatus()).isEqualTo(DocumentVersionStatus.EMBEDDING);
+
+ ArgumentCaptor> embeddingsCaptor = ArgumentCaptor.forClass(List.class);
+ then(embeddingRepository).should().saveAllAndFlush(embeddingsCaptor.capture());
+ Embedding first = embeddingsCaptor.getValue().get(0);
+ assertThat(first.getChunk()).isSameAs(chunks.get(0));
+ assertThat(first.getDocumentVersion()).isSameAs(documentVersion);
+ assertThat(first.getEmbeddingModel()).isSameAs(embeddingModel);
+ assertThat(first.getDimension()).isEqualTo(EmbeddingModelFixture.DIMENSION);
+ assertThat(first.getVector()).hasSize(EmbeddingModelFixture.DIMENSION);
+ assertThat(first.getStatus()).isEqualTo(EmbeddingStatus.ACTIVE);
+ then(indexingEventRepository).shouldHaveNoInteractions();
+ }
+
+ @Test
+ @DisplayName("완료 단계에서 선행 요청의 전체 저장을 발견하면 Insert 없이 재생한다")
+ void complete_replaysConcurrentWinner() {
+ prepareEntities(DocumentVersionStatus.EMBEDDING);
+ List chunks = List.of(chunk(0, "첫 번째"), chunk(1, "두 번째"));
+ givenValidContext(chunks);
+ given(embeddingRepository.countByDocumentVersionIdAndEmbeddingModelId(VERSION_ID, MODEL_ID))
+ .willReturn(2L);
+
+ CompletionResult result = service.complete(
+ JOB_ID,
+ ATTEMPT_ID,
+ WORKER_ID,
+ CLAIM_TOKEN,
+ work(chunks),
+ List.of(draft(chunks.get(0), 0.1f), draft(chunks.get(1), 0.2f))
+ );
+
+ assertThat(result.created()).isFalse();
+ assertThat(result.embeddingCount()).isEqualTo(2);
+ then(embeddingRepository).should(never()).saveAllAndFlush(any());
+ }
+
+ @Test
+ @DisplayName("외부 호출 중 Chunk Text가 바뀌면 Draft 저장을 거부한다")
+ void complete_rejectsChangedChunk() {
+ prepareEntities(DocumentVersionStatus.EMBEDDING);
+ List originalChunks = List.of(chunk(0, "원본"));
+ EmbeddingWork work = work(originalChunks);
+ List changedChunks = List.of(chunk(0, "변경"));
+ givenValidContext(changedChunks);
+
+ assertThatThrownBy(() -> service.complete(
+ JOB_ID,
+ ATTEMPT_ID,
+ WORKER_ID,
+ CLAIM_TOKEN,
+ work,
+ List.of(draft(originalChunks.get(0), 0.1f))
+ ))
+ .isInstanceOfSatisfying(DocGridException.class,
+ exception -> assertThat(exception.getErrorCode())
+ .isEqualTo(ErrorCode.DOCUMENT_CHUNKS_INCONSISTENT));
+
+ then(embeddingRepository).should(never()).saveAllAndFlush(any());
+ }
+
+ @Test
+ @DisplayName("Vector 내용과 Hash가 일치하지 않으면 Draft 저장을 거부한다")
+ void complete_rejectsTamperedVectorHash() {
+ prepareEntities(DocumentVersionStatus.EMBEDDING);
+ List chunks = List.of(chunk(0, "본문"));
+ givenValidContext(chunks);
+ float[] vector = vector(0.1f);
+ DocumentEmbeddingDraft tampered = new DocumentEmbeddingDraft(
+ chunks.get(0).getId(),
+ 0,
+ CONTENT_HASH,
+ vector,
+ "0".repeat(64)
+ );
+
+ assertThatThrownBy(() -> service.complete(
+ JOB_ID,
+ ATTEMPT_ID,
+ WORKER_ID,
+ CLAIM_TOKEN,
+ work(chunks),
+ List.of(tampered)
+ ))
+ .isInstanceOfSatisfying(DocGridException.class,
+ exception -> assertThat(exception.getErrorCode())
+ .isEqualTo(ErrorCode.EMBEDDING_VECTOR_INVALID));
+
+ then(embeddingRepository).should(never()).saveAllAndFlush(any());
+ }
+
+ @Test
+ @DisplayName("준비 Snapshot의 Model이 현재 Job 고정 Model과 다르면 Chunk를 다시 읽지 않는다")
+ void complete_rejectsChangedModel() {
+ prepareEntities(DocumentVersionStatus.EMBEDDING);
+ List chunks = List.of(chunk(0, "본문"));
+ given(embeddingJobRepository.findByIdForUpdate(JOB_ID)).willReturn(Optional.of(embeddingJob));
+ given(embeddingJobAttemptRepository.findByEmbeddingJobIdAndClaimToken(JOB_ID, CLAIM_TOKEN))
+ .willReturn(Optional.of(attempt));
+ given(documentVersionRepository.findByIdForUpdate(VERSION_ID)).willReturn(Optional.of(documentVersion));
+ EmbeddingWork changedModelWork = new EmbeddingWork(
+ VERSION_ID,
+ 99L,
+ EmbeddingModelFixture.DIMENSION,
+ work(chunks).chunks()
+ );
+
+ assertThatThrownBy(() -> service.complete(
+ JOB_ID,
+ ATTEMPT_ID,
+ WORKER_ID,
+ CLAIM_TOKEN,
+ changedModelWork,
+ List.of(draft(chunks.get(0), 0.1f))
+ ))
+ .isInstanceOfSatisfying(DocGridException.class,
+ exception -> assertThat(exception.getErrorCode())
+ .isEqualTo(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT));
+
+ then(documentChunkRepository).shouldHaveNoInteractions();
+ }
+
+ private void givenValidContext(List chunks) {
+ given(embeddingJobRepository.findByIdForUpdate(JOB_ID)).willReturn(Optional.of(embeddingJob));
+ given(embeddingJobAttemptRepository.findByEmbeddingJobIdAndClaimToken(JOB_ID, CLAIM_TOKEN))
+ .willReturn(Optional.of(attempt));
+ given(documentVersionRepository.findByIdForUpdate(VERSION_ID)).willReturn(Optional.of(documentVersion));
+ given(documentChunkRepository.findAllByDocumentVersionIdOrderByChunkIndexAsc(VERSION_ID))
+ .willReturn(chunks);
+ }
+
+ private void prepareEntities(DocumentVersionStatus versionStatus) {
+ WorkerNode worker = WorkerNode.builder()
+ .workerName("worker")
+ .instanceId("instance")
+ .status(WorkerStatus.ACTIVE)
+ .startedAt(NOW.minusHours(1))
+ .build();
+ ReflectionTestUtils.setField(worker, "id", WORKER_ID);
+
+ Document document = Document.builder()
+ .title("문서")
+ .documentType(DocumentType.TXT)
+ .sourceType(DocumentSourceType.UPLOAD)
+ .status(DocumentStatus.INDEXING)
+ .visibility(VisibilityType.PRIVATE)
+ .build();
+ ReflectionTestUtils.setField(document, "id", 3L);
+
+ documentVersion = DocumentVersion.builder()
+ .document(document)
+ .versionNo(1)
+ .status(versionStatus)
+ .build();
+ ReflectionTestUtils.setField(documentVersion, "id", VERSION_ID);
+
+ embeddingModel = EmbeddingModelFixture.createDefaultModel();
+ ReflectionTestUtils.setField(embeddingModel, "id", MODEL_ID);
+
+ embeddingJob = EmbeddingJob.builder()
+ .documentVersion(documentVersion)
+ .embeddingModel(embeddingModel)
+ .status(EmbeddingJobStatus.PROCESSING)
+ .maxRetryCount(3)
+ .build();
+ ReflectionTestUtils.setField(embeddingJob, "id", JOB_ID);
+
+ attempt = EmbeddingJobAttempt.builder()
+ .embeddingJob(embeddingJob)
+ .workerNode(worker)
+ .attemptNo(1)
+ .claimToken(CLAIM_TOKEN)
+ .status(AttemptStatus.STARTED)
+ .startedAt(NOW.minusMinutes(1))
+ .build();
+ ReflectionTestUtils.setField(attempt, "id", ATTEMPT_ID);
+ }
+
+ private DocumentChunk chunk(int index, String text) {
+ DocumentChunk chunk = DocumentChunk.builder()
+ .documentVersion(documentVersion)
+ .chunkIndex(index)
+ .chunkText(text)
+ .tokenCount(1)
+ .charStart(index * 10)
+ .charEnd(index * 10 + text.length())
+ .contentHash(CONTENT_HASH)
+ .build();
+ ReflectionTestUtils.setField(chunk, "id", 20L + index);
+ return chunk;
+ }
+
+ private EmbeddingWork work(List chunks) {
+ return new EmbeddingWork(
+ VERSION_ID,
+ MODEL_ID,
+ EmbeddingModelFixture.DIMENSION,
+ chunks.stream()
+ .map(chunk -> new ChunkSnapshot(
+ chunk.getId(),
+ chunk.getChunkIndex(),
+ chunk.getChunkText(),
+ chunk.getContentHash()
+ ))
+ .toList()
+ );
+ }
+
+ private DocumentEmbeddingDraft draft(DocumentChunk chunk, float firstValue) {
+ float[] vector = vector(firstValue);
+ return new DocumentEmbeddingDraft(
+ chunk.getId(),
+ chunk.getChunkIndex(),
+ chunk.getContentHash(),
+ vector,
+ EmbeddingVectorSupport.calculateHash(vector)
+ );
+ }
+
+ private float[] vector(float firstValue) {
+ float[] vector = new float[EmbeddingModelFixture.DIMENSION];
+ vector[0] = firstValue;
+ return vector;
+ }
+}
diff --git a/src/test/java/com/opensource/docgrid/domain/embedding/service/query/QueryEmbeddingServiceTest.java b/src/test/java/com/opensource/docgrid/domain/embedding/service/query/QueryEmbeddingServiceTest.java
index 61975e2..0cc6410 100644
--- a/src/test/java/com/opensource/docgrid/domain/embedding/service/query/QueryEmbeddingServiceTest.java
+++ b/src/test/java/com/opensource/docgrid/domain/embedding/service/query/QueryEmbeddingServiceTest.java
@@ -3,55 +3,37 @@
import static org.assertj.core.api.Assertions.assertThat;
import static org.assertj.core.api.Assertions.assertThatThrownBy;
import static org.mockito.BDDMockito.given;
-import static org.mockito.Mockito.doReturn;
-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.Answers;
+import org.mockito.InjectMocks;
import org.mockito.Mock;
import org.mockito.junit.jupiter.MockitoExtension;
-import org.mockito.junit.jupiter.MockitoSettings;
-import org.mockito.quality.Strictness;
-import org.springframework.web.client.ResourceAccessException;
-import org.springframework.web.client.RestClient;
+import com.opensource.docgrid.domain.embedding.client.EmbeddingClient;
import com.opensource.docgrid.domain.embedding.dto.EmbedResult;
-import com.opensource.docgrid.domain.embedding.dto.response.EmbedServerResponse;
import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel;
import com.opensource.docgrid.domain.embedding.fixture.EmbeddingModelFixture;
import com.opensource.docgrid.global.exception.DocGridException;
import com.opensource.docgrid.global.exception.ErrorCode;
@ExtendWith(MockitoExtension.class)
-@MockitoSettings(strictness = Strictness.LENIENT)
@DisplayName("QueryEmbeddingService 단위 테스트")
class QueryEmbeddingServiceTest {
@Mock private EmbeddingModelQueryService embeddingModelQueryService;
- @Mock private RestClient restClient;
- @Mock(answer = Answers.RETURNS_SELF) private RestClient.RequestBodyUriSpec requestBodyUriSpec;
- @Mock private RestClient.ResponseSpec responseSpec;
-
- private QueryEmbeddingService queryEmbeddingService;
-
- @BeforeEach
- void setUp() {
- queryEmbeddingService = new QueryEmbeddingService(embeddingModelQueryService, restClient);
- doReturn(requestBodyUriSpec).when(restClient).post();
- doReturn(responseSpec).when(requestBodyUriSpec).retrieve();
- }
+ @Mock private EmbeddingClient embeddingClient;
+ @InjectMocks private QueryEmbeddingService queryEmbeddingService;
@Test
@DisplayName("정상 케이스: 텍스트를 1024차원 벡터로 변환한다")
void embed_success() {
EmbeddingModel model = EmbeddingModelFixture.createDefaultModel();
float[] vector = new float[EmbeddingModelFixture.DIMENSION];
- EmbedServerResponse serverResponse = new EmbedServerResponse(vector);
given(embeddingModelQueryService.getActiveModel()).willReturn(model);
- given(responseSpec.body(EmbedServerResponse.class)).willReturn(serverResponse);
+ given(embeddingClient.embed("검색어")).willReturn(vector);
EmbedResult result = queryEmbeddingService.embed("검색어");
@@ -63,10 +45,9 @@ void embed_success() {
@DisplayName("차원 불일치: 응답 차원이 모델 차원과 다르면 EMBEDDING_DIMENSION_MISMATCH 예외가 발생한다")
void embed_dimensionMismatch_throwsException() {
EmbeddingModel model = EmbeddingModelFixture.createDefaultModel();
- EmbedServerResponse serverResponse = new EmbedServerResponse(new float[768]);
given(embeddingModelQueryService.getActiveModel()).willReturn(model);
- given(responseSpec.body(EmbedServerResponse.class)).willReturn(serverResponse);
+ given(embeddingClient.embed("검색어")).willReturn(new float[768]);
assertThatThrownBy(() -> queryEmbeddingService.embed("검색어"))
.isInstanceOf(DocGridException.class)
@@ -74,16 +55,15 @@ void embed_dimensionMismatch_throwsException() {
}
@Test
- @DisplayName("서버 장애: RestClientException 발생 시 EMBEDDING_SERVER_UNAVAILABLE 예외가 발생한다")
- void embed_serverUnavailable_throwsException() {
+ @DisplayName("빈 응답: Vector가 null이면 EMBEDDING_DIMENSION_MISMATCH 예외가 발생한다")
+ void embed_nullVector_throwsException() {
EmbeddingModel model = EmbeddingModelFixture.createDefaultModel();
given(embeddingModelQueryService.getActiveModel()).willReturn(model);
- given(responseSpec.body(EmbedServerResponse.class))
- .willThrow(new ResourceAccessException("Connection refused"));
+ given(embeddingClient.embed("검색어")).willReturn(null);
assertThatThrownBy(() -> queryEmbeddingService.embed("검색어"))
.isInstanceOf(DocGridException.class)
- .hasFieldOrPropertyWithValue("errorCode", ErrorCode.EMBEDDING_SERVER_UNAVAILABLE);
+ .hasFieldOrPropertyWithValue("errorCode", ErrorCode.EMBEDDING_DIMENSION_MISMATCH);
}
}