diff --git a/backend/src/main/java/com/opensource/docgrid/domain/document/entity/DocumentVersion.java b/backend/src/main/java/com/opensource/docgrid/domain/document/entity/DocumentVersion.java
index 1db0798..99b4686 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/document/entity/DocumentVersion.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/document/entity/DocumentVersion.java
@@ -155,6 +155,19 @@ public void markIndexed(LocalDateTime indexedAt) {
this.indexedAt = indexedAt;
}
+ /**
+ * 현재 검색 Version의 Vector 손상이 확인됐을 때 기존 Chunk Set부터 다시 임베딩하도록 되돌린다.
+ *
+ *
호출 Service가 현재 Version·Chunk 존재·Job 부재를 잠금 상태에서 검증해야 한다.
+ */
+ public void reopenIndexedForVectorRepair() {
+ if (status != DocumentVersionStatus.INDEXED) {
+ throw new IllegalStateException("INDEXED 상태의 문서 버전만 Vector 복구를 시작할 수 있습니다.");
+ }
+ status = DocumentVersionStatus.CHUNKED;
+ indexedAt = null;
+ }
+
/**
* 최종 실패한 Version을 수동 재처리가 다시 진행할 수 있는 재개 지점으로 되돌린다.
*
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/document/repository/DocumentChunkRepository.java b/backend/src/main/java/com/opensource/docgrid/domain/document/repository/DocumentChunkRepository.java
index 10158f9..12c32ab 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/document/repository/DocumentChunkRepository.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/document/repository/DocumentChunkRepository.java
@@ -3,6 +3,7 @@
import java.util.List;
import org.springframework.data.jpa.repository.JpaRepository;
+import org.springframework.data.jpa.repository.Query;
import com.opensource.docgrid.domain.document.entity.DocumentChunk;
@@ -19,5 +20,17 @@ public interface DocumentChunkRepository extends JpaRepository findAllByDocumentVersionIdOrderByChunkIndexAsc(Long documentVersionId);
}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/document/repository/DocumentVersionRepository.java b/backend/src/main/java/com/opensource/docgrid/domain/document/repository/DocumentVersionRepository.java
index 7dff659..b5d3c33 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/document/repository/DocumentVersionRepository.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/document/repository/DocumentVersionRepository.java
@@ -1,6 +1,7 @@
package com.opensource.docgrid.domain.document.repository;
import java.util.Collection;
+import java.util.List;
import java.util.Optional;
import jakarta.persistence.LockModeType;
@@ -9,6 +10,7 @@
import org.springframework.data.jpa.repository.Lock;
import org.springframework.data.jpa.repository.Query;
import org.springframework.data.repository.query.Param;
+import org.springframework.data.domain.Pageable;
import com.opensource.docgrid.domain.document.entity.DocumentVersion;
import com.opensource.docgrid.domain.document.enums.DocumentVersionStatus;
@@ -21,6 +23,22 @@
*/
public interface DocumentVersionRepository extends JpaRepository {
+ /**
+ * Reconciler가 전체 Version을 Offset 없이 작은 ID Cursor Batch로 순회한다.
+ */
+ @Query("""
+ SELECT version
+ FROM DocumentVersion version
+ JOIN FETCH version.document document
+ LEFT JOIN FETCH document.currentVersion
+ WHERE version.id > :cursor
+ ORDER BY version.id ASC
+ """)
+ List findReconciliationBatchAfterId(
+ @Param("cursor") Long cursor,
+ Pageable pageable
+ );
+
boolean existsByDocumentIdAndStatusIn(Long documentId, Collection statuses);
Optional findTopByDocumentIdOrderByVersionNoDesc(Long documentId);
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/embedding/repository/EmbeddingJobRepository.java b/backend/src/main/java/com/opensource/docgrid/domain/embedding/repository/EmbeddingJobRepository.java
index 1b6d06e..31396b3 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/embedding/repository/EmbeddingJobRepository.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/embedding/repository/EmbeddingJobRepository.java
@@ -33,6 +33,11 @@ Optional findTopByDocumentVersionIdAndEmbeddingModelIdOrderByIdDes
Long embeddingModelId
);
+ boolean existsByDocumentVersionIdAndStatusIn(
+ Long documentVersionId,
+ Collection statuses
+ );
+
/**
* 관리자 목록 화면에 필요한 연관관계를 함께 조회하면서 선택 필터와 Pagination을 적용한다.
*
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/embedding/repository/EmbeddingRepository.java b/backend/src/main/java/com/opensource/docgrid/domain/embedding/repository/EmbeddingRepository.java
index 971024a..f81d94f 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/embedding/repository/EmbeddingRepository.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/embedding/repository/EmbeddingRepository.java
@@ -1,5 +1,7 @@
package com.opensource.docgrid.domain.embedding.repository;
+import java.util.List;
+
import org.springframework.data.jpa.repository.JpaRepository;
import org.springframework.data.jpa.repository.Modifying;
import org.springframework.data.jpa.repository.Query;
@@ -31,6 +33,36 @@ long countByDocumentVersionIdAndEmbeddingModelIdAndStatus(
EmbeddingStatus status
);
+ @Query("""
+ SELECT embedding.chunk.id
+ FROM Embedding embedding
+ WHERE embedding.documentVersion.id = :documentVersionId
+ AND embedding.embeddingModel.id = :embeddingModelId
+ """)
+ List findChunkIdsByDocumentVersionIdAndEmbeddingModelId(
+ @Param("documentVersionId") Long documentVersionId,
+ @Param("embeddingModelId") Long embeddingModelId
+ );
+
+ long countByDocumentIdAndStatus(Long documentId, EmbeddingStatus status);
+
+ /**
+ * Chunk·Document·Version·Model 원장 중 하나라도 사라진 Embedding 행을 탐지한다.
+ */
+ @Query(value = """
+ SELECT COUNT(*)
+ FROM embeddings embedding
+ LEFT JOIN document_chunks chunk ON chunk.id = embedding.chunk_id
+ LEFT JOIN documents document ON document.id = embedding.document_id
+ LEFT JOIN document_versions version ON version.id = embedding.document_version_id
+ LEFT JOIN embedding_models model ON model.id = embedding.embedding_model_id
+ WHERE chunk.id IS NULL
+ OR document.id IS NULL
+ OR version.id IS NULL
+ OR model.id IS NULL
+ """, nativeQuery = true)
+ long countOrphanedRows();
+
/**
* 인덱싱 완료 시 이전 현재 Version 또는 최종 실패 대상의 검색 가능한 Embedding을 한 SQL로 비활성화한다.
*
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionService.java b/backend/src/main/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionService.java
index c7ff459..62a2308 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionService.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionService.java
@@ -5,6 +5,7 @@
import java.util.ArrayList;
import java.util.List;
import java.util.Objects;
+import java.util.Set;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;
@@ -39,6 +40,7 @@
*
* 두 단계는 Job을 먼저, Version을 다음 순서로 잠가 Claim 교체와 같은 Version의 동시 실행을
* 직렬화한다. 준비 단계는 Job에 고정된 Model과 Chunk 불변 Snapshot만 외부 호출 구간에 전달한다.
+ * Reconciler 복구 Job은 기존 Vector를 보존하고 같은 Version·Model에서 실제 누락된 Chunk만 채운다.
*/
@Service
@RequiredArgsConstructor
@@ -150,7 +152,7 @@ public CompletionResult complete(
validateChunks(documentVersion, chunks);
validatePreparedChunks(preparedWork, chunks);
- // 4. 동시 요청이 먼저 전체 저장했으면 기존 결과를 재생하고 부분 저장은 내부 모순으로 거부한다.
+ // 4. 동시 요청이 먼저 전체 저장했으면 재생하고 Reconciler 복구의 부분 Set은 누락분만 채운다.
EmbeddingState state = resolveState(documentVersion, embeddingModel, chunks.size());
if (state == EmbeddingState.REPLAY) {
return result(jobId, attemptId, documentVersion, embeddingModel, chunks.size(), false);
@@ -159,16 +161,18 @@ public CompletionResult complete(
throw new DocGridException(ErrorCode.DOCUMENT_VERSION_EMBEDDING_NOT_ALLOWED);
}
- // 5. 모든 Draft를 다시 검증하고 같은 Version·Model의 ACTIVE Embedding Set으로 원자 저장한다.
+ // 5. 모든 Draft를 다시 검증하고 아직 없는 Chunk의 ACTIVE Embedding만 원자 저장한다.
List embeddings = toEntities(
documentVersion,
embeddingModel,
chunks,
drafts
);
- embeddingRepository.saveAllAndFlush(embeddings);
+ if (!embeddings.isEmpty()) {
+ embeddingRepository.saveAllAndFlush(embeddings);
+ }
- return result(jobId, attemptId, documentVersion, embeddingModel, embeddings.size(), true);
+ return result(jobId, attemptId, documentVersion, embeddingModel, chunks.size(), true);
}
private EmbeddingJob findLockedJob(Long jobId) {
@@ -282,14 +286,17 @@ private EmbeddingState resolveState(
embeddingModel.getId()
);
+ if (embeddingCount > chunkCount) {
+ throw new DocGridException(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT);
+ }
if (documentVersion.getStatus() == DocumentVersionStatus.CHUNKED) {
- if (embeddingCount != 0) {
- throw new DocGridException(ErrorCode.DOCUMENT_EMBEDDINGS_INCONSISTENT);
+ if (embeddingCount == chunkCount) {
+ return EmbeddingState.REPLAY;
}
return EmbeddingState.WORK;
}
if (documentVersion.getStatus() == DocumentVersionStatus.EMBEDDING) {
- if (embeddingCount == 0) {
+ if (embeddingCount < chunkCount) {
return EmbeddingState.WORK;
}
if (embeddingCount == chunkCount) {
@@ -317,10 +324,20 @@ private List toEntities(
}
List embeddings = new ArrayList<>(drafts.size());
+ Set existingChunkIds = Set.copyOf(
+ embeddingRepository.findChunkIdsByDocumentVersionIdAndEmbeddingModelId(
+ documentVersion.getId(),
+ embeddingModel.getId()
+ )
+ );
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());
+ if (existingChunkIds.contains(chunk.getId())) {
+ // 기존 Vector는 보존하고 누락 Chunk만 채워 자동복구가 물리 삭제를 요구하지 않게 한다.
+ continue;
+ }
embeddings.add(Embedding.builder()
.chunk(chunk)
.document(documentVersion.getDocument())
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/embedding/service/command/DocumentIndexingCompletionService.java b/backend/src/main/java/com/opensource/docgrid/domain/embedding/service/command/DocumentIndexingCompletionService.java
index 437502e..19bd396 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/embedding/service/command/DocumentIndexingCompletionService.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/embedding/service/command/DocumentIndexingCompletionService.java
@@ -410,7 +410,6 @@ private CompletionCounts validateEmbeddingSet(
if (activeJobCount != 1
|| chunkCount <= 0
- || allEmbeddingCount != chunkCount
|| modelEmbeddingCount != chunkCount
|| activeEmbeddingCount != chunkCount
|| invalidEmbeddingCount != 0) {
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/embedding/service/command/IndexedVersionVectorRepairService.java b/backend/src/main/java/com/opensource/docgrid/domain/embedding/service/command/IndexedVersionVectorRepairService.java
new file mode 100644
index 0000000..87f4f6c
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/embedding/service/command/IndexedVersionVectorRepairService.java
@@ -0,0 +1,114 @@
+package com.opensource.docgrid.domain.embedding.service.command;
+
+import java.util.EnumSet;
+import java.util.Objects;
+import java.util.Set;
+import java.util.UUID;
+
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+
+import com.opensource.docgrid.domain.document.entity.Document;
+import com.opensource.docgrid.domain.document.entity.DocumentVersion;
+import com.opensource.docgrid.domain.document.enums.DocumentStatus;
+import com.opensource.docgrid.domain.document.enums.DocumentVersionStatus;
+import com.opensource.docgrid.domain.document.repository.DocumentChunkRepository;
+import com.opensource.docgrid.domain.document.repository.DocumentRepository;
+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.enums.EmbeddingJobStatus;
+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.global.exception.DocGridException;
+import com.opensource.docgrid.global.exception.ErrorCode;
+
+import lombok.RequiredArgsConstructor;
+
+/**
+ * 현재 INDEXED Version을 기존 Vector 보존 상태로 Chunk 기반 재임베딩 Queue에 되돌린다.
+ *
+ * Version·Document 잠금 안에서 현재 검색 대상과 부분 Set을 검증하고 새 Job만 만든다. Worker는 기존
+ * Vector를 덮거나 지우지 않고 누락 Chunk만 채우며, 다른 모델의 Vector도 모델 전환 이력으로 보존한다.
+ */
+@Service
+@RequiredArgsConstructor
+@Transactional
+public class IndexedVersionVectorRepairService {
+
+ private static final int DEFAULT_JOB_PRIORITY = 0;
+ private static final int MAX_RETRY_COUNT = 3;
+ private static final Set LIVE_JOB_STATUSES = EnumSet.of(
+ EmbeddingJobStatus.PENDING,
+ EmbeddingJobStatus.PROCESSING
+ );
+
+ private final DocumentVersionRepository documentVersionRepository;
+ private final DocumentRepository documentRepository;
+ private final DocumentChunkRepository documentChunkRepository;
+ private final EmbeddingJobRepository embeddingJobRepository;
+ private final EmbeddingRepository embeddingRepository;
+
+ public EmbeddingJob repair(
+ Long documentVersionId,
+ EmbeddingModel embeddingModel,
+ UUID sourceEventId
+ ) {
+ // 1. Worker 완료 경로와 같은 Version → Document 순서로 잠그고 현재 검색 대상인지 검증한다.
+ DocumentVersion version = documentVersionRepository.findByIdForUpdate(documentVersionId)
+ .orElseThrow(() -> new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT));
+ Document document = documentRepository.findByIdForUpdate(version.getDocument().getId())
+ .orElseThrow(() -> new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT));
+ validateTarget(document, version, embeddingModel);
+
+ // 2. 검색 노출을 중단하고 기존 Chunk Set을 재사용하는 상태로 원자 전환한다.
+ version.reopenIndexedForVectorRepair();
+ document.markIndexing();
+
+ // 3. Repair Event를 원인으로 가진 Job을 만들어 Event 재전달에서도 한 건만 유지한다.
+ return embeddingJobRepository.save(
+ EmbeddingJob.builder()
+ .documentVersion(version)
+ .embeddingModel(embeddingModel)
+ .sourceEventId(sourceEventId)
+ .status(EmbeddingJobStatus.PENDING)
+ .priority(DEFAULT_JOB_PRIORITY)
+ .maxRetryCount(MAX_RETRY_COUNT)
+ .build()
+ );
+ }
+
+ private void validateTarget(
+ Document document,
+ DocumentVersion version,
+ EmbeddingModel embeddingModel
+ ) {
+ if (version.getStatus() != DocumentVersionStatus.INDEXED
+ || document.getStatus() != DocumentStatus.INDEXED
+ || document.getCurrentVersion() == null
+ || !Objects.equals(document.getCurrentVersion().getId(), version.getId())
+ || !documentChunkRepository.existsByDocumentVersionId(version.getId())
+ || embeddingJobRepository.existsByDocumentVersionIdAndStatusIn(
+ version.getId(), LIVE_JOB_STATUSES
+ )) {
+ throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT);
+ }
+ long chunkCount = documentChunkRepository.countByDocumentVersionId(version.getId());
+ long allModelEmbeddingCount = embeddingRepository.countByDocumentVersionIdAndEmbeddingModelId(
+ version.getId(),
+ embeddingModel.getId()
+ );
+ long activeEmbeddingCount = embeddingRepository
+ .countByDocumentVersionIdAndEmbeddingModelIdAndStatus(
+ version.getId(),
+ embeddingModel.getId(),
+ EmbeddingStatus.ACTIVE
+ );
+ if (allModelEmbeddingCount >= chunkCount
+ || activeEmbeddingCount >= chunkCount
+ || allModelEmbeddingCount != activeEmbeddingCount) {
+ throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT);
+ }
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/config/SyncReconciliationProperties.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/config/SyncReconciliationProperties.java
new file mode 100644
index 0000000..0d46446
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/config/SyncReconciliationProperties.java
@@ -0,0 +1,49 @@
+package com.opensource.docgrid.domain.sync.config;
+
+import java.time.Duration;
+
+import org.springframework.boot.context.properties.ConfigurationProperties;
+import org.springframework.stereotype.Component;
+import org.springframework.validation.annotation.Validated;
+
+import com.opensource.docgrid.domain.sync.enums.SyncReconciliationMode;
+
+import jakarta.validation.constraints.AssertTrue;
+import jakarta.validation.constraints.Min;
+import jakarta.validation.constraints.NotNull;
+import lombok.Getter;
+import lombok.Setter;
+
+/**
+ * Reconciler의 활성화 여부, Batch 크기, 실행 주기와 정체 판단 임계시간을 바인딩한다.
+ */
+@Getter
+@Setter
+@Validated
+@Component
+@ConfigurationProperties(prefix = "sync.reconciliation")
+public class SyncReconciliationProperties {
+
+ private boolean enabled = false;
+
+ @NotNull
+ private SyncReconciliationMode mode = SyncReconciliationMode.DRY_RUN;
+
+ @Min(1)
+ private int batchSize = 100;
+
+ @NotNull
+ private Duration interval = Duration.ofMinutes(5);
+
+ @NotNull
+ private Duration stalledThreshold = Duration.ofMinutes(15);
+
+ @AssertTrue(message = "Reconciliation 실행 주기와 정체 임계시간은 0보다 커야 합니다.")
+ public boolean isTimingValid() {
+ return isPositive(interval) && isPositive(stalledThreshold);
+ }
+
+ private boolean isPositive(Duration duration) {
+ return duration != null && !duration.isZero() && !duration.isNegative();
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/config/SyncReconciliationSchedulingConfig.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/config/SyncReconciliationSchedulingConfig.java
new file mode 100644
index 0000000..b934f6c
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/config/SyncReconciliationSchedulingConfig.java
@@ -0,0 +1,14 @@
+package com.opensource.docgrid.domain.sync.config;
+
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.scheduling.annotation.EnableScheduling;
+
+/**
+ * Reconciliation이 활성화된 실행 인스턴스에서 Cursor Scheduler를 켠다.
+ */
+@Configuration
+@EnableScheduling
+@ConditionalOnProperty(prefix = "sync.reconciliation", name = "enabled", havingValue = "true")
+public class SyncReconciliationSchedulingConfig {
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/dto/SyncReconciliationBatchResult.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/dto/SyncReconciliationBatchResult.java
new file mode 100644
index 0000000..1c2652a
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/dto/SyncReconciliationBatchResult.java
@@ -0,0 +1,17 @@
+package com.opensource.docgrid.domain.sync.dto;
+
+import java.util.UUID;
+
+/**
+ * 한 Reconciliation Cursor Batch의 실행 범위와 탐지·복구 요청 집계를 전달한다.
+ */
+public record SyncReconciliationBatchResult(
+ UUID runId,
+ long startCursor,
+ long endCursor,
+ int scannedCount,
+ int detectedCount,
+ int repairRequestedCount,
+ boolean hasMore
+) {
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/entity/SyncConsistencyIssue.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/entity/SyncConsistencyIssue.java
new file mode 100644
index 0000000..fbd2272
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/entity/SyncConsistencyIssue.java
@@ -0,0 +1,177 @@
+package com.opensource.docgrid.domain.sync.entity;
+
+import java.time.LocalDateTime;
+import java.util.UUID;
+
+import com.opensource.docgrid.domain.document.entity.Document;
+import com.opensource.docgrid.domain.document.entity.DocumentVersion;
+import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencyIssueStatus;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencyIssueType;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencySeverity;
+import com.opensource.docgrid.global.common.entity.BaseEntity;
+
+import jakarta.persistence.Column;
+import jakarta.persistence.Entity;
+import jakarta.persistence.EnumType;
+import jakarta.persistence.Enumerated;
+import jakarta.persistence.FetchType;
+import jakarta.persistence.GeneratedValue;
+import jakarta.persistence.GenerationType;
+import jakarta.persistence.Id;
+import jakarta.persistence.Index;
+import jakarta.persistence.JoinColumn;
+import jakarta.persistence.ManyToOne;
+import jakarta.persistence.Table;
+import jakarta.persistence.UniqueConstraint;
+import lombok.AccessLevel;
+import lombok.Builder;
+import lombok.Getter;
+import lombok.NoArgsConstructor;
+
+/**
+ * Reconciler가 발견한 동일한 데이터 불일치를 하나의 관리 가능한 Issue로 보존한다.
+ *
+ * issueKey는 반복 검사에서도 같은 행으로 수렴시키며, 기대·실제 Snapshot과 복구 Event를 함께 남긴다.
+ * 이 Entity는 문제를 설명하고 복구 생명주기를 추적하지만 currentVersion 변경이나 물리 삭제를 직접
+ * 수행하지 않는다.
+ */
+@Getter
+@Entity
+@NoArgsConstructor(access = AccessLevel.PROTECTED)
+@Table(
+ name = "sync_consistency_issues",
+ uniqueConstraints = @UniqueConstraint(
+ name = "uk_sync_consistency_issues_issue_key",
+ columnNames = "issue_key"
+ ),
+ indexes = {
+ @Index(name = "idx_sync_consistency_issues_status_last_detected", columnList = "status, last_detected_at, id"),
+ @Index(name = "idx_sync_consistency_issues_document_version", columnList = "document_version_id, status"),
+ @Index(name = "idx_sync_consistency_issues_repair_event", columnList = "repair_event_id")
+ }
+)
+public class SyncConsistencyIssue extends BaseEntity {
+
+ @Id
+ @GeneratedValue(strategy = GenerationType.IDENTITY)
+ private Long id;
+
+ @Column(name = "issue_key", nullable = false, updatable = false, length = 255)
+ private String issueKey;
+
+ @Enumerated(EnumType.STRING)
+ @Column(name = "issue_type", nullable = false, updatable = false, length = 50)
+ private SyncConsistencyIssueType issueType;
+
+ @Enumerated(EnumType.STRING)
+ @Column(nullable = false, length = 20)
+ private SyncConsistencySeverity severity;
+
+ @Enumerated(EnumType.STRING)
+ @Column(nullable = false, length = 20)
+ private SyncConsistencyIssueStatus status;
+
+ @ManyToOne(fetch = FetchType.LAZY)
+ @JoinColumn(name = "document_id")
+ private Document document;
+
+ @ManyToOne(fetch = FetchType.LAZY)
+ @JoinColumn(name = "document_version_id")
+ private DocumentVersion documentVersion;
+
+ @ManyToOne(fetch = FetchType.LAZY)
+ @JoinColumn(name = "embedding_model_id")
+ private EmbeddingModel embeddingModel;
+
+ @Column(name = "expected_json", columnDefinition = "TEXT")
+ private String expectedJson;
+
+ @Column(name = "actual_json", columnDefinition = "TEXT")
+ private String actualJson;
+
+ @Column(name = "detected_at", nullable = false, updatable = false)
+ private LocalDateTime detectedAt;
+
+ @Column(name = "last_detected_at", nullable = false)
+ private LocalDateTime lastDetectedAt;
+
+ @Column(name = "repair_event_id")
+ private UUID repairEventId;
+
+ @Column(name = "repair_attempt_count", nullable = false)
+ private int repairAttemptCount;
+
+ @Column(name = "resolved_at")
+ private LocalDateTime resolvedAt;
+
+ @Column(name = "resolution_message", columnDefinition = "TEXT")
+ private String resolutionMessage;
+
+ @Builder
+ public SyncConsistencyIssue(
+ String issueKey,
+ SyncConsistencyIssueType issueType,
+ SyncConsistencySeverity severity,
+ Document document,
+ DocumentVersion documentVersion,
+ EmbeddingModel embeddingModel,
+ String expectedJson,
+ String actualJson,
+ LocalDateTime detectedAt
+ ) {
+ this.issueKey = issueKey;
+ this.issueType = issueType;
+ this.severity = severity;
+ this.status = SyncConsistencyIssueStatus.OPEN;
+ this.document = document;
+ this.documentVersion = documentVersion;
+ this.embeddingModel = embeddingModel;
+ this.expectedJson = expectedJson;
+ this.actualJson = actualJson;
+ this.detectedAt = detectedAt;
+ this.lastDetectedAt = detectedAt;
+ }
+
+ public void detectAgain(
+ SyncConsistencySeverity newSeverity,
+ String newExpectedJson,
+ String newActualJson,
+ LocalDateTime detectedAgainAt
+ ) {
+ severity = newSeverity;
+ expectedJson = newExpectedJson;
+ actualJson = newActualJson;
+ lastDetectedAt = detectedAgainAt;
+ if (status == SyncConsistencyIssueStatus.RESOLVED) {
+ status = SyncConsistencyIssueStatus.OPEN;
+ resolvedAt = null;
+ resolutionMessage = null;
+ }
+ }
+
+ public void markRepairing(UUID newRepairEventId, LocalDateTime requestedAt) {
+ if (status == SyncConsistencyIssueStatus.IGNORED || newRepairEventId == null || requestedAt == null) {
+ throw new IllegalStateException("무시되지 않은 Issue에 유효한 Repair Event가 필요합니다.");
+ }
+ status = SyncConsistencyIssueStatus.REPAIRING;
+ repairEventId = newRepairEventId;
+ repairAttemptCount++;
+ lastDetectedAt = requestedAt;
+ }
+
+ public void resolve(LocalDateTime resolvedAt, String message) {
+ if (status == SyncConsistencyIssueStatus.IGNORED || resolvedAt == null) {
+ throw new IllegalStateException("무시된 Issue는 자동 해결할 수 없습니다.");
+ }
+ status = SyncConsistencyIssueStatus.RESOLVED;
+ this.resolvedAt = resolvedAt;
+ this.resolutionMessage = message;
+ }
+
+ public void ignore(LocalDateTime ignoredAt, String message) {
+ status = SyncConsistencyIssueStatus.IGNORED;
+ resolvedAt = ignoredAt;
+ resolutionMessage = message;
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/entity/SyncReconciliationRun.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/entity/SyncReconciliationRun.java
new file mode 100644
index 0000000..c1cc62d
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/entity/SyncReconciliationRun.java
@@ -0,0 +1,118 @@
+package com.opensource.docgrid.domain.sync.entity;
+
+import java.time.LocalDateTime;
+import java.util.UUID;
+
+import com.opensource.docgrid.domain.sync.enums.SyncReconciliationMode;
+import com.opensource.docgrid.domain.sync.enums.SyncReconciliationStatus;
+import com.opensource.docgrid.global.common.entity.BaseEntity;
+
+import jakarta.persistence.Column;
+import jakarta.persistence.Entity;
+import jakarta.persistence.EnumType;
+import jakarta.persistence.Enumerated;
+import jakarta.persistence.GeneratedValue;
+import jakarta.persistence.GenerationType;
+import jakarta.persistence.Id;
+import jakarta.persistence.Index;
+import jakarta.persistence.Table;
+import jakarta.persistence.UniqueConstraint;
+import lombok.AccessLevel;
+import lombok.Builder;
+import lombok.Getter;
+import lombok.NoArgsConstructor;
+
+/**
+ * Reconciliation Batch의 범위와 탐지·복구 결과를 운영 지표로 보존한다.
+ *
+ *
문서별 Issue 내용은 SyncConsistencyIssue가 담당하고, 이 Entity는 한 번의 실행 단위만 요약한다.
+ */
+@Getter
+@Entity
+@NoArgsConstructor(access = AccessLevel.PROTECTED)
+@Table(
+ name = "sync_reconciliation_runs",
+ uniqueConstraints = @UniqueConstraint(name = "uk_sync_reconciliation_runs_run_id", columnNames = "run_id"),
+ indexes = {
+ @Index(name = "idx_sync_reconciliation_runs_started_at", columnList = "started_at, id"),
+ @Index(name = "idx_sync_reconciliation_runs_status", columnList = "status, started_at")
+ }
+)
+public class SyncReconciliationRun extends BaseEntity {
+
+ @Id
+ @GeneratedValue(strategy = GenerationType.IDENTITY)
+ private Long id;
+
+ @Column(name = "run_id", nullable = false, updatable = false)
+ private UUID runId;
+
+ @Enumerated(EnumType.STRING)
+ @Column(nullable = false, updatable = false, length = 20)
+ private SyncReconciliationMode mode;
+
+ @Enumerated(EnumType.STRING)
+ @Column(nullable = false, length = 20)
+ private SyncReconciliationStatus status;
+
+ @Column(name = "start_cursor", nullable = false, updatable = false)
+ private long startCursor;
+
+ @Column(name = "end_cursor", nullable = false)
+ private long endCursor;
+
+ @Column(name = "scanned_count", nullable = false)
+ private int scannedCount;
+
+ @Column(name = "detected_count", nullable = false)
+ private int detectedCount;
+
+ @Column(name = "repair_requested_count", nullable = false)
+ private int repairRequestedCount;
+
+ @Column(name = "started_at", nullable = false, updatable = false)
+ private LocalDateTime startedAt;
+
+ @Column(name = "completed_at")
+ private LocalDateTime completedAt;
+
+ @Column(name = "error_code", length = 100)
+ private String errorCode;
+
+ @Builder
+ public SyncReconciliationRun(
+ UUID runId,
+ SyncReconciliationMode mode,
+ long startCursor,
+ LocalDateTime startedAt
+ ) {
+ this.runId = runId;
+ this.mode = mode;
+ this.status = SyncReconciliationStatus.RUNNING;
+ this.startCursor = startCursor;
+ this.endCursor = startCursor;
+ this.startedAt = startedAt;
+ }
+
+ public void complete(
+ long completedCursor,
+ int scanned,
+ int detected,
+ int repairRequested,
+ LocalDateTime completedAt
+ ) {
+ status = SyncReconciliationStatus.COMPLETED;
+ endCursor = completedCursor;
+ scannedCount = scanned;
+ detectedCount = detected;
+ repairRequestedCount = repairRequested;
+ this.completedAt = completedAt;
+ errorCode = null;
+ }
+
+ public void fail(String errorCode, LocalDateTime failedAt) {
+ status = SyncReconciliationStatus.FAILED;
+ this.errorCode = errorCode;
+ completedAt = failedAt;
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncConsistencyIssueStatus.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncConsistencyIssueStatus.java
new file mode 100644
index 0000000..08993af
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncConsistencyIssueStatus.java
@@ -0,0 +1,11 @@
+package com.opensource.docgrid.domain.sync.enums;
+
+/**
+ * 일관성 Issue의 탐지, 자동 복구, 해결과 관리자 무시 생명주기를 정의한다.
+ */
+public enum SyncConsistencyIssueStatus {
+ OPEN,
+ REPAIRING,
+ RESOLVED,
+ IGNORED
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncConsistencyIssueType.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncConsistencyIssueType.java
new file mode 100644
index 0000000..8c96042
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncConsistencyIssueType.java
@@ -0,0 +1,15 @@
+package com.opensource.docgrid.domain.sync.enums;
+
+/**
+ * 문서 원장과 Job·Chunk·Vector 파생 데이터 사이에서 Reconciler가 탐지하는 불일치 유형이다.
+ */
+public enum SyncConsistencyIssueType {
+ MISSING_JOB,
+ MISSING_CHUNKS,
+ MISSING_EMBEDDINGS,
+ MODEL_MISMATCH,
+ INVALID_CURRENT_VERSION,
+ STALLED_VERSION,
+ DELETED_DOCUMENT_RESIDUE,
+ ORPHANED_DATA
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncConsistencySeverity.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncConsistencySeverity.java
new file mode 100644
index 0000000..2925e81
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncConsistencySeverity.java
@@ -0,0 +1,10 @@
+package com.opensource.docgrid.domain.sync.enums;
+
+/**
+ * 동기화 불일치가 검색 가용성과 데이터 신뢰도에 미치는 영향을 구분한다.
+ */
+public enum SyncConsistencySeverity {
+ WARNING,
+ ERROR,
+ CRITICAL
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncReconciliationMode.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncReconciliationMode.java
new file mode 100644
index 0000000..7b9ffc3
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncReconciliationMode.java
@@ -0,0 +1,9 @@
+package com.opensource.docgrid.domain.sync.enums;
+
+/**
+ * Reconciliation Batch가 탐지만 수행하는지 안전한 복구 Event까지 생성하는지 구분한다.
+ */
+public enum SyncReconciliationMode {
+ DRY_RUN,
+ REPAIR
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncReconciliationStatus.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncReconciliationStatus.java
new file mode 100644
index 0000000..9e93135
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncReconciliationStatus.java
@@ -0,0 +1,10 @@
+package com.opensource.docgrid.domain.sync.enums;
+
+/**
+ * Reconciliation Batch 실행 이력의 시작과 종결 상태를 나타낸다.
+ */
+public enum SyncReconciliationStatus {
+ RUNNING,
+ COMPLETED,
+ FAILED
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/lifecycle/SyncReconciliationScheduler.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/lifecycle/SyncReconciliationScheduler.java
new file mode 100644
index 0000000..c42475e
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/lifecycle/SyncReconciliationScheduler.java
@@ -0,0 +1,57 @@
+package com.opensource.docgrid.domain.sync.lifecycle;
+
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicLong;
+
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.scheduling.annotation.Scheduled;
+import org.springframework.stereotype.Component;
+
+import com.opensource.docgrid.domain.sync.config.SyncReconciliationProperties;
+import com.opensource.docgrid.domain.sync.dto.SyncReconciliationBatchResult;
+import com.opensource.docgrid.domain.sync.service.SyncReconciliationOrchestrator;
+
+import lombok.RequiredArgsConstructor;
+import lombok.extern.slf4j.Slf4j;
+
+/**
+ * 전체 DocumentVersion을 작은 ID Cursor Batch로 순회하고 마지막 Batch 뒤 처음부터 재검사한다.
+ */
+@Slf4j
+@Component
+@RequiredArgsConstructor
+@ConditionalOnProperty(prefix = "sync.reconciliation", name = "enabled", havingValue = "true")
+public class SyncReconciliationScheduler {
+
+ private final SyncReconciliationOrchestrator syncReconciliationOrchestrator;
+ private final SyncReconciliationProperties properties;
+ private final AtomicBoolean reconciling = new AtomicBoolean(false);
+ private final AtomicLong cursor = new AtomicLong(0L);
+
+ @Scheduled(
+ fixedDelayString = "${sync.reconciliation.interval:5m}",
+ initialDelayString = "${sync.reconciliation.interval:5m}"
+ )
+ public void reconcile() {
+ if (!reconciling.compareAndSet(false, true)) {
+ return;
+ }
+ try {
+ long startCursor = cursor.get();
+ SyncReconciliationBatchResult result = syncReconciliationOrchestrator.reconcileBatch(
+ startCursor,
+ properties.getMode()
+ );
+ cursor.set(result.hasMore() ? result.endCursor() : 0L);
+ } catch (RuntimeException exception) {
+ // 실패한 Cursor를 유지해 다음 주기에 같은 범위를 다시 검사한다.
+ log.error(
+ "Sync Reconciliation에 실패했습니다. cursor={}, errorType={}",
+ cursor.get(),
+ exception.getClass().getSimpleName()
+ );
+ } finally {
+ reconciling.set(false);
+ }
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/repository/SyncConsistencyIssueRepository.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/repository/SyncConsistencyIssueRepository.java
new file mode 100644
index 0000000..32181dd
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/repository/SyncConsistencyIssueRepository.java
@@ -0,0 +1,26 @@
+package com.opensource.docgrid.domain.sync.repository;
+
+import java.util.List;
+import java.util.Optional;
+
+import org.springframework.data.jpa.repository.JpaRepository;
+
+import com.opensource.docgrid.domain.sync.entity.SyncConsistencyIssue;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencyIssueStatus;
+
+/**
+ * 일관성 Issue의 멱등 Key 조회와 관리자 목록 기반을 제공한다.
+ */
+public interface SyncConsistencyIssueRepository extends JpaRepository {
+
+ Optional findByIssueKey(String issueKey);
+
+ List findAllByDocumentVersionIdAndStatusIn(
+ Long documentVersionId,
+ List statuses
+ );
+
+ List findAllByDocumentVersionIsNullAndStatusIn(
+ List statuses
+ );
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/repository/SyncReconciliationRunRepository.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/repository/SyncReconciliationRunRepository.java
new file mode 100644
index 0000000..2d48690
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/repository/SyncReconciliationRunRepository.java
@@ -0,0 +1,16 @@
+package com.opensource.docgrid.domain.sync.repository;
+
+import java.util.Optional;
+import java.util.UUID;
+
+import org.springframework.data.jpa.repository.JpaRepository;
+
+import com.opensource.docgrid.domain.sync.entity.SyncReconciliationRun;
+
+/**
+ * Reconciliation Batch 실행 이력을 영속화한다.
+ */
+public interface SyncReconciliationRunRepository extends JpaRepository {
+
+ Optional findByRunId(UUID runId);
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncConsistencyInspector.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncConsistencyInspector.java
new file mode 100644
index 0000000..ca26a28
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncConsistencyInspector.java
@@ -0,0 +1,247 @@
+package com.opensource.docgrid.domain.sync.service;
+
+import java.time.LocalDateTime;
+import java.util.ArrayList;
+import java.util.EnumSet;
+import java.util.List;
+import java.util.Objects;
+import java.util.Set;
+
+import org.springframework.stereotype.Component;
+
+import com.opensource.docgrid.domain.document.entity.Document;
+import com.opensource.docgrid.domain.document.entity.DocumentVersion;
+import com.opensource.docgrid.domain.document.enums.DocumentStatus;
+import com.opensource.docgrid.domain.document.enums.DocumentVersionStatus;
+import com.opensource.docgrid.domain.document.repository.DocumentChunkRepository;
+import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel;
+import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus;
+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.sync.config.SyncReconciliationProperties;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencyIssueType;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencySeverity;
+
+import lombok.RequiredArgsConstructor;
+
+/**
+ * 한 DocumentVersion의 원장 상태와 Job·Chunk·Vector 파생 상태를 읽기 전용으로 비교한다.
+ *
+ * 검사는 데이터 변경을 수행하지 않으며, 안전한 자동복구 가능 여부를 Observation에 표시한다.
+ * currentVersion 모순, 기존 Vector 행의 상태 손상과 삭제 잔여물은 항상 보고 전용이다.
+ */
+@Component
+@RequiredArgsConstructor
+public class SyncConsistencyInspector {
+
+ private static final Set PROCESSING_STATUSES = EnumSet.of(
+ DocumentVersionStatus.UPLOADED,
+ DocumentVersionStatus.PARSING,
+ DocumentVersionStatus.CHUNKED,
+ DocumentVersionStatus.EMBEDDING
+ );
+ private static final Set LIVE_JOB_STATUSES = EnumSet.of(
+ EmbeddingJobStatus.PENDING,
+ EmbeddingJobStatus.PROCESSING
+ );
+
+ private final EmbeddingJobRepository embeddingJobRepository;
+ private final DocumentChunkRepository documentChunkRepository;
+ private final EmbeddingRepository embeddingRepository;
+ private final SyncReconciliationProperties properties;
+
+ public List inspect(
+ DocumentVersion version,
+ EmbeddingModel activeModel,
+ LocalDateTime inspectedAt
+ ) {
+ Document document = version.getDocument();
+ List observations = new ArrayList<>();
+ long chunkCount = documentChunkRepository.countByDocumentVersionId(version.getId());
+ long activeModelEmbeddingCount = embeddingRepository
+ .countByDocumentVersionIdAndEmbeddingModelIdAndStatus(
+ version.getId(), activeModel.getId(), EmbeddingStatus.ACTIVE
+ );
+ boolean liveJobExists = embeddingJobRepository.existsByDocumentVersionIdAndStatusIn(
+ version.getId(), LIVE_JOB_STATUSES
+ );
+
+ inspectMissingJob(version, activeModel, liveJobExists, observations);
+ inspectChunksAndEmbeddings(version, activeModel, chunkCount, activeModelEmbeddingCount, observations);
+ inspectCurrentVersion(document, version, activeModel, observations);
+ inspectStalledVersion(version, activeModel, liveJobExists, inspectedAt, observations);
+ inspectDeletedResidue(document, version, activeModel, observations);
+ return observations;
+ }
+
+ private void inspectMissingJob(
+ DocumentVersion version,
+ EmbeddingModel activeModel,
+ boolean liveJobExists,
+ List observations
+ ) {
+ if (PROCESSING_STATUSES.contains(version.getStatus()) && !liveJobExists) {
+ observations.add(observation(
+ SyncConsistencyIssueType.MISSING_JOB,
+ SyncConsistencySeverity.ERROR,
+ version,
+ activeModel,
+ "{\"liveJob\":true}",
+ "{\"liveJob\":false,\"versionStatus\":\"%s\"}".formatted(version.getStatus()),
+ true
+ ));
+ }
+ }
+
+ private void inspectChunksAndEmbeddings(
+ DocumentVersion version,
+ EmbeddingModel activeModel,
+ long chunkCount,
+ long embeddingCount,
+ List observations
+ ) {
+ if ((version.getStatus() == DocumentVersionStatus.CHUNKED
+ || version.getStatus() == DocumentVersionStatus.EMBEDDING
+ || version.getStatus() == DocumentVersionStatus.INDEXED)
+ && chunkCount == 0) {
+ observations.add(observation(
+ SyncConsistencyIssueType.MISSING_CHUNKS,
+ SyncConsistencySeverity.CRITICAL,
+ version,
+ activeModel,
+ "{\"chunkCount\":\">0\"}",
+ "{\"chunkCount\":0}",
+ false
+ ));
+ }
+ boolean currentVersion = version.getDocument().getCurrentVersion() != null
+ && Objects.equals(version.getDocument().getCurrentVersion().getId(), version.getId());
+ if (version.getStatus() == DocumentVersionStatus.INDEXED
+ && currentVersion
+ && chunkCount > embeddingCount) {
+ long currentModelCount = embeddingRepository.countByDocumentVersionIdAndEmbeddingModelId(
+ version.getId(),
+ activeModel.getId()
+ );
+ long allModelCount = embeddingRepository.countByDocumentVersionId(version.getId());
+ SyncConsistencyIssueType issueType = currentModelCount == 0 && allModelCount >= chunkCount
+ ? SyncConsistencyIssueType.MODEL_MISMATCH
+ : SyncConsistencyIssueType.MISSING_EMBEDDINGS;
+ observations.add(observation(
+ issueType,
+ SyncConsistencySeverity.CRITICAL,
+ version,
+ activeModel,
+ "{\"activeEmbeddingCount\":%d}".formatted(chunkCount),
+ "{\"activeEmbeddingCount\":%d,\"currentModelEmbeddingCount\":%d,"
+ + "\"allModelEmbeddingCount\":%d}"
+ .formatted(embeddingCount, currentModelCount, allModelCount),
+ currentModelCount == embeddingCount && currentModelCount < chunkCount
+ ));
+ }
+ }
+
+ private void inspectCurrentVersion(
+ Document document,
+ DocumentVersion version,
+ EmbeddingModel activeModel,
+ List observations
+ ) {
+ if (document.getStatus() != DocumentStatus.INDEXED) {
+ return;
+ }
+ DocumentVersion current = document.getCurrentVersion();
+ boolean missingCurrentVersion = current == null;
+ boolean inspectingCurrentVersion = current != null && Objects.equals(current.getId(), version.getId());
+ if (missingCurrentVersion
+ || (inspectingCurrentVersion
+ && (version.getStatus() != DocumentVersionStatus.INDEXED
+ || current.getDocument() == null
+ || !Objects.equals(current.getDocument().getId(), document.getId())))) {
+ observations.add(observation(
+ SyncConsistencyIssueType.INVALID_CURRENT_VERSION,
+ SyncConsistencySeverity.CRITICAL,
+ version,
+ activeModel,
+ "{\"currentVersionStatus\":\"INDEXED\",\"sameDocument\":true}",
+ current == null
+ ? "{\"currentVersion\":null}"
+ : "{\"currentVersionStatus\":\"%s\"}".formatted(version.getStatus()),
+ false
+ ));
+ }
+ }
+
+ private void inspectStalledVersion(
+ DocumentVersion version,
+ EmbeddingModel activeModel,
+ boolean liveJobExists,
+ LocalDateTime inspectedAt,
+ List observations
+ ) {
+ if (!PROCESSING_STATUSES.contains(version.getStatus())
+ || version.getUpdatedAt() == null
+ || !version.getUpdatedAt().isBefore(inspectedAt.minus(properties.getStalledThreshold()))) {
+ return;
+ }
+ observations.add(observation(
+ SyncConsistencyIssueType.STALLED_VERSION,
+ liveJobExists ? SyncConsistencySeverity.WARNING : SyncConsistencySeverity.ERROR,
+ version,
+ activeModel,
+ "{\"updatedAfter\":\"%s\"}".formatted(inspectedAt.minus(properties.getStalledThreshold())),
+ "{\"updatedAt\":\"%s\",\"liveJob\":%s}"
+ .formatted(version.getUpdatedAt(), liveJobExists),
+ false
+ ));
+ }
+
+ private void inspectDeletedResidue(
+ Document document,
+ DocumentVersion version,
+ EmbeddingModel activeModel,
+ List observations
+ ) {
+ if (document.getStatus() == DocumentStatus.DELETED
+ && document.getCurrentVersion() != null
+ && Objects.equals(document.getCurrentVersion().getId(), version.getId())
+ && embeddingRepository.countByDocumentIdAndStatus(document.getId(), EmbeddingStatus.ACTIVE) > 0) {
+ observations.add(observation(
+ SyncConsistencyIssueType.DELETED_DOCUMENT_RESIDUE,
+ SyncConsistencySeverity.CRITICAL,
+ version,
+ activeModel,
+ "{\"activeEmbeddingCount\":0}",
+ "{\"activeEmbeddingCount\":\">0\"}",
+ false
+ ));
+ }
+ }
+
+ private SyncConsistencyObservation observation(
+ SyncConsistencyIssueType type,
+ SyncConsistencySeverity severity,
+ DocumentVersion version,
+ EmbeddingModel model,
+ String expectedJson,
+ String actualJson,
+ boolean repairable
+ ) {
+ String issueKey = type == SyncConsistencyIssueType.INVALID_CURRENT_VERSION
+ || type == SyncConsistencyIssueType.DELETED_DOCUMENT_RESIDUE
+ ? "%s:DOCUMENT:%d".formatted(type, version.getDocument().getId())
+ : "%s:VERSION:%d:MODEL:%d".formatted(type, version.getId(), model.getId());
+ return new SyncConsistencyObservation(
+ issueKey,
+ type,
+ severity,
+ version.getDocument(),
+ version,
+ model,
+ expectedJson,
+ actualJson,
+ repairable
+ );
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncConsistencyObservation.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncConsistencyObservation.java
new file mode 100644
index 0000000..826b870
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncConsistencyObservation.java
@@ -0,0 +1,26 @@
+package com.opensource.docgrid.domain.sync.service;
+
+import com.opensource.docgrid.domain.document.entity.Document;
+import com.opensource.docgrid.domain.document.entity.DocumentVersion;
+import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencyIssueType;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencySeverity;
+
+/**
+ * 한 번의 Version 검사에서 발견한 기대 상태와 실제 상태의 불일치 Snapshot이다.
+ *
+ * Issue Service가 이 값을 멱등 Issue로 영속화하고, repairable 값이 true인 경우에만 Outbox 기반
+ * 재인덱싱을 요청한다.
+ */
+public record SyncConsistencyObservation(
+ String issueKey,
+ SyncConsistencyIssueType issueType,
+ SyncConsistencySeverity severity,
+ Document document,
+ DocumentVersion documentVersion,
+ EmbeddingModel embeddingModel,
+ String expectedJson,
+ String actualJson,
+ boolean repairable
+) {
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncOrphanInspector.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncOrphanInspector.java
new file mode 100644
index 0000000..b84489c
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncOrphanInspector.java
@@ -0,0 +1,43 @@
+package com.opensource.docgrid.domain.sync.service;
+
+import java.util.List;
+
+import org.springframework.stereotype.Component;
+
+import com.opensource.docgrid.domain.document.repository.DocumentChunkRepository;
+import com.opensource.docgrid.domain.embedding.repository.EmbeddingRepository;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencyIssueType;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencySeverity;
+
+import lombok.RequiredArgsConstructor;
+
+/**
+ * FK 우회나 부분 복구로 원장 관계를 잃은 Chunk·Embedding 행을 전역 읽기 검사로 탐지한다.
+ */
+@Component
+@RequiredArgsConstructor
+public class SyncOrphanInspector {
+
+ private final DocumentChunkRepository documentChunkRepository;
+ private final EmbeddingRepository embeddingRepository;
+
+ public List inspect() {
+ long orphanChunkCount = documentChunkRepository.countOrphanedRows();
+ long orphanEmbeddingCount = embeddingRepository.countOrphanedRows();
+ if (orphanChunkCount == 0 && orphanEmbeddingCount == 0) {
+ return List.of();
+ }
+ return List.of(new SyncConsistencyObservation(
+ "ORPHANED_DATA:GLOBAL",
+ SyncConsistencyIssueType.ORPHANED_DATA,
+ SyncConsistencySeverity.CRITICAL,
+ null,
+ null,
+ null,
+ "{\"orphanChunkCount\":0,\"orphanEmbeddingCount\":0}",
+ "{\"orphanChunkCount\":%d,\"orphanEmbeddingCount\":%d}"
+ .formatted(orphanChunkCount, orphanEmbeddingCount),
+ false
+ ));
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncReconciliationOrchestrator.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncReconciliationOrchestrator.java
new file mode 100644
index 0000000..2d373fa
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncReconciliationOrchestrator.java
@@ -0,0 +1,50 @@
+package com.opensource.docgrid.domain.sync.service;
+
+import java.util.UUID;
+
+import org.springframework.stereotype.Service;
+
+import com.opensource.docgrid.domain.sync.dto.SyncReconciliationBatchResult;
+import com.opensource.docgrid.domain.sync.enums.SyncReconciliationMode;
+import com.opensource.docgrid.domain.sync.service.command.SyncReconciliationBatchService;
+import com.opensource.docgrid.domain.sync.service.command.SyncReconciliationRunService;
+import com.opensource.docgrid.global.exception.DocGridException;
+
+import lombok.RequiredArgsConstructor;
+
+/**
+ * Reconciliation 실행 이력과 검사 Transaction을 조율하고 실패 결과를 독립적으로 보존한다.
+ */
+@Service
+@RequiredArgsConstructor
+public class SyncReconciliationOrchestrator {
+
+ private final SyncReconciliationRunService syncReconciliationRunService;
+ private final SyncReconciliationBatchService syncReconciliationBatchService;
+
+ public SyncReconciliationBatchResult reconcileBatch(
+ long startCursor,
+ SyncReconciliationMode mode
+ ) {
+ UUID runId = syncReconciliationRunService.start(mode, startCursor);
+ try {
+ SyncReconciliationBatchResult result = syncReconciliationBatchService.reconcile(
+ runId,
+ startCursor,
+ mode
+ );
+ syncReconciliationRunService.complete(result);
+ return result;
+ } catch (RuntimeException exception) {
+ syncReconciliationRunService.fail(runId, diagnosticCode(exception));
+ throw exception;
+ }
+ }
+
+ private String diagnosticCode(RuntimeException exception) {
+ if (exception instanceof DocGridException docGridException) {
+ return docGridException.getErrorCode().getCode();
+ }
+ return exception.getClass().getSimpleName();
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncConsistencyIssueService.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncConsistencyIssueService.java
new file mode 100644
index 0000000..1988696
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncConsistencyIssueService.java
@@ -0,0 +1,91 @@
+package com.opensource.docgrid.domain.sync.service.command;
+
+import java.time.LocalDateTime;
+import java.util.List;
+import java.util.Set;
+
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+
+import com.opensource.docgrid.domain.sync.entity.SyncConsistencyIssue;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencyIssueStatus;
+import com.opensource.docgrid.domain.sync.repository.SyncConsistencyIssueRepository;
+import com.opensource.docgrid.domain.sync.service.SyncConsistencyObservation;
+
+import lombok.RequiredArgsConstructor;
+
+/**
+ * 반복 탐지된 같은 불일치를 하나의 Issue로 수렴시키고 정상화된 Issue를 해결한다.
+ *
+ * 자동 복구 Event 생성은 Reconciler가 담당하며 이 Service는 Issue 생명주기만 변경한다.
+ */
+@Service
+@RequiredArgsConstructor
+@Transactional
+public class SyncConsistencyIssueService {
+
+ private static final List ACTIVE_STATUSES = List.of(
+ SyncConsistencyIssueStatus.OPEN,
+ SyncConsistencyIssueStatus.REPAIRING
+ );
+
+ private final SyncConsistencyIssueRepository syncConsistencyIssueRepository;
+
+ public SyncConsistencyIssue detect(
+ SyncConsistencyObservation observation,
+ LocalDateTime detectedAt
+ ) {
+ return syncConsistencyIssueRepository.findByIssueKey(observation.issueKey())
+ .map(issue -> detectAgain(issue, observation, detectedAt))
+ .orElseGet(() -> syncConsistencyIssueRepository.save(
+ SyncConsistencyIssue.builder()
+ .issueKey(observation.issueKey())
+ .issueType(observation.issueType())
+ .severity(observation.severity())
+ .document(observation.document())
+ .documentVersion(observation.documentVersion())
+ .embeddingModel(observation.embeddingModel())
+ .expectedJson(observation.expectedJson())
+ .actualJson(observation.actualJson())
+ .detectedAt(detectedAt)
+ .build()
+ ));
+ }
+
+ public void resolveMissingObservations(
+ Long documentVersionId,
+ Set detectedIssueKeys,
+ LocalDateTime resolvedAt
+ ) {
+ syncConsistencyIssueRepository
+ .findAllByDocumentVersionIdAndStatusIn(documentVersionId, ACTIVE_STATUSES)
+ .stream()
+ .filter(issue -> !detectedIssueKeys.contains(issue.getIssueKey()))
+ .forEach(issue -> issue.resolve(resolvedAt, "Reconciliation 재검사에서 정상 상태를 확인했습니다."));
+ }
+
+ public void resolveMissingGlobalObservations(
+ Set detectedIssueKeys,
+ LocalDateTime resolvedAt
+ ) {
+ syncConsistencyIssueRepository
+ .findAllByDocumentVersionIsNullAndStatusIn(ACTIVE_STATUSES)
+ .stream()
+ .filter(issue -> !detectedIssueKeys.contains(issue.getIssueKey()))
+ .forEach(issue -> issue.resolve(resolvedAt, "Reconciliation 전역 재검사에서 정상 상태를 확인했습니다."));
+ }
+
+ private SyncConsistencyIssue detectAgain(
+ SyncConsistencyIssue issue,
+ SyncConsistencyObservation observation,
+ LocalDateTime detectedAt
+ ) {
+ issue.detectAgain(
+ observation.severity(),
+ observation.expectedJson(),
+ observation.actualJson(),
+ detectedAt
+ );
+ return issue;
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncReconciliationBatchService.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncReconciliationBatchService.java
new file mode 100644
index 0000000..08b1901
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncReconciliationBatchService.java
@@ -0,0 +1,130 @@
+package com.opensource.docgrid.domain.sync.service.command;
+
+import java.time.Clock;
+import java.time.LocalDateTime;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+import java.util.UUID;
+
+import org.springframework.data.domain.PageRequest;
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+
+import com.opensource.docgrid.domain.document.entity.DocumentVersion;
+import com.opensource.docgrid.domain.document.repository.DocumentVersionRepository;
+import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel;
+import com.opensource.docgrid.domain.embedding.service.query.EmbeddingModelQueryService;
+import com.opensource.docgrid.domain.sync.config.SyncReconciliationProperties;
+import com.opensource.docgrid.domain.sync.dto.SyncReconciliationBatchResult;
+import com.opensource.docgrid.domain.sync.entity.SyncConsistencyIssue;
+import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencyIssueStatus;
+import com.opensource.docgrid.domain.sync.enums.SyncReconciliationMode;
+import com.opensource.docgrid.domain.sync.service.SyncConsistencyInspector;
+import com.opensource.docgrid.domain.sync.service.SyncConsistencyObservation;
+import com.opensource.docgrid.domain.sync.service.SyncOrphanInspector;
+
+import lombok.RequiredArgsConstructor;
+
+/**
+ * 작은 ID Cursor Batch에서 원장과 파생 데이터를 비교하고 Issue 및 안전한 Repair Event를 원자적으로 기록한다.
+ *
+ * 누락 Job·Vector처럼 재실행 가능한 불일치만 Outbox로 복구하며, 데이터 삭제나 currentVersion 변경은 수행하지 않는다.
+ */
+@Service
+@RequiredArgsConstructor
+@Transactional
+public class SyncReconciliationBatchService {
+
+ private final DocumentVersionRepository documentVersionRepository;
+ private final EmbeddingModelQueryService embeddingModelQueryService;
+ private final SyncConsistencyInspector syncConsistencyInspector;
+ private final SyncConsistencyIssueService syncConsistencyIssueService;
+ private final SyncOrphanInspector syncOrphanInspector;
+ private final SyncEventWriter syncEventWriter;
+ private final SyncReconciliationProperties properties;
+ private final Clock clock;
+
+ public SyncReconciliationBatchResult reconcile(
+ UUID runId,
+ long startCursor,
+ SyncReconciliationMode mode
+ ) {
+ LocalDateTime inspectedAt = LocalDateTime.now(clock);
+ EmbeddingModel activeModel = embeddingModelQueryService.getActiveModel();
+ List versions = documentVersionRepository.findReconciliationBatchAfterId(
+ startCursor,
+ PageRequest.of(0, properties.getBatchSize())
+ );
+ int detectedCount = 0;
+ int repairRequestedCount = 0;
+
+ // 1. 첫 Cursor에서는 Version FK 밖의 전역 고아 데이터도 한 번 검사한다.
+ if (startCursor == 0L) {
+ Set orphanIssueKeys = new HashSet<>();
+ for (SyncConsistencyObservation observation : syncOrphanInspector.inspect()) {
+ orphanIssueKeys.add(observation.issueKey());
+ syncConsistencyIssueService.detect(observation, inspectedAt);
+ detectedCount++;
+ }
+ syncConsistencyIssueService.resolveMissingGlobalObservations(orphanIssueKeys, inspectedAt);
+ }
+
+ // 2. 각 Version을 독립 Issue Key 집합으로 검사해 반복 실행을 같은 Issue 행에 수렴시킨다.
+ for (DocumentVersion version : versions) {
+ List observations = syncConsistencyInspector.inspect(
+ version,
+ activeModel,
+ inspectedAt
+ );
+ Set detectedIssueKeys = new HashSet<>();
+ for (SyncConsistencyObservation observation : observations) {
+ detectedIssueKeys.add(observation.issueKey());
+ SyncConsistencyIssue issue = syncConsistencyIssueService.detect(observation, inspectedAt);
+ detectedCount++;
+
+ // 3. REPAIR 모드에서도 명시적으로 안전하다고 판정된 OPEN Issue만 Outbox에 기록한다.
+ if (mode == SyncReconciliationMode.REPAIR
+ && observation.repairable()
+ && issue.getStatus() == SyncConsistencyIssueStatus.OPEN) {
+ SyncOutboxEvent repairEvent = syncEventWriter.recordDocumentReindexRequested(
+ version,
+ activeModel,
+ repairRequestKey(issue)
+ );
+ issue.markRepairing(repairEvent.getEventId(), inspectedAt);
+ repairRequestedCount++;
+ }
+ }
+
+ // 4. 이전 실행의 활성 Issue가 이번 검사에서 사라졌다면 정상화된 것으로 종결한다.
+ syncConsistencyIssueService.resolveMissingObservations(
+ version.getId(),
+ detectedIssueKeys,
+ inspectedAt
+ );
+ }
+
+ // 5. 마지막 ID를 다음 Cursor로 반환해 다음 Batch가 Offset 재탐색 없이 이어지게 한다.
+ long endCursor = versions.isEmpty()
+ ? startCursor
+ : versions.get(versions.size() - 1).getId();
+ return new SyncReconciliationBatchResult(
+ runId,
+ startCursor,
+ endCursor,
+ versions.size(),
+ detectedCount,
+ repairRequestedCount,
+ versions.size() == properties.getBatchSize()
+ );
+ }
+
+ private String repairRequestKey(SyncConsistencyIssue issue) {
+ return "reconcile:%s:attempt:%d".formatted(
+ issue.getIssueKey(),
+ issue.getRepairAttemptCount() + 1
+ );
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncReconciliationRunService.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncReconciliationRunService.java
new file mode 100644
index 0000000..73c17af
--- /dev/null
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncReconciliationRunService.java
@@ -0,0 +1,63 @@
+package com.opensource.docgrid.domain.sync.service.command;
+
+import java.time.Clock;
+import java.time.LocalDateTime;
+import java.util.UUID;
+
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Propagation;
+import org.springframework.transaction.annotation.Transactional;
+
+import com.opensource.docgrid.domain.sync.dto.SyncReconciliationBatchResult;
+import com.opensource.docgrid.domain.sync.entity.SyncReconciliationRun;
+import com.opensource.docgrid.domain.sync.enums.SyncReconciliationMode;
+import com.opensource.docgrid.domain.sync.repository.SyncReconciliationRunRepository;
+
+import lombok.RequiredArgsConstructor;
+
+/**
+ * Reconciliation 실행 이력을 검사 Transaction과 분리해 시작·완료·실패 상태로 보존한다.
+ *
+ * 검사 Transaction이 Rollback돼도 FAILED 실행 이력은 독립 Transaction으로 남는다.
+ */
+@Service
+@RequiredArgsConstructor
+@Transactional(propagation = Propagation.REQUIRES_NEW)
+public class SyncReconciliationRunService {
+
+ private final SyncReconciliationRunRepository syncReconciliationRunRepository;
+ private final Clock clock;
+
+ public UUID start(SyncReconciliationMode mode, long startCursor) {
+ UUID runId = UUID.randomUUID();
+ syncReconciliationRunRepository.save(
+ SyncReconciliationRun.builder()
+ .runId(runId)
+ .mode(mode)
+ .startCursor(startCursor)
+ .startedAt(LocalDateTime.now(clock))
+ .build()
+ );
+ return runId;
+ }
+
+ public void complete(SyncReconciliationBatchResult result) {
+ SyncReconciliationRun run = getRun(result.runId());
+ run.complete(
+ result.endCursor(),
+ result.scannedCount(),
+ result.detectedCount(),
+ result.repairRequestedCount(),
+ LocalDateTime.now(clock)
+ );
+ }
+
+ public void fail(UUID runId, String errorCode) {
+ getRun(runId).fail(errorCode, LocalDateTime.now(clock));
+ }
+
+ private SyncReconciliationRun getRun(UUID runId) {
+ return syncReconciliationRunRepository.findByRunId(runId)
+ .orElseThrow(() -> new IllegalStateException("Reconciliation 실행 이력을 찾을 수 없습니다."));
+ }
+}
diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/DocumentVersionSyncEventHandler.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/DocumentVersionSyncEventHandler.java
index 00fde35..cc52689 100644
--- a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/DocumentVersionSyncEventHandler.java
+++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/DocumentVersionSyncEventHandler.java
@@ -16,6 +16,7 @@
import com.opensource.docgrid.domain.embedding.repository.EmbeddingJobRepository;
import com.opensource.docgrid.domain.embedding.repository.EmbeddingModelRepository;
import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobManualRetryService;
+import com.opensource.docgrid.domain.embedding.service.command.IndexedVersionVectorRepairService;
import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent;
import com.opensource.docgrid.domain.sync.enums.SyncEventType;
import com.opensource.docgrid.domain.sync.service.SyncEventHandler;
@@ -41,11 +42,18 @@ public class DocumentVersionSyncEventHandler implements SyncEventHandler {
SyncEventType.DOCUMENT_VERSION_CREATED,
SyncEventType.DOCUMENT_REINDEX_REQUESTED
);
+ private static final Set RECOVERABLE_PROCESSING_STATUSES = EnumSet.of(
+ DocumentVersionStatus.UPLOADED,
+ DocumentVersionStatus.PARSING,
+ DocumentVersionStatus.CHUNKED,
+ DocumentVersionStatus.EMBEDDING
+ );
private final DocumentVersionRepository documentVersionRepository;
private final EmbeddingModelRepository embeddingModelRepository;
private final EmbeddingJobRepository embeddingJobRepository;
private final EmbeddingJobManualRetryService embeddingJobManualRetryService;
+ private final IndexedVersionVectorRepairService indexedVersionVectorRepairService;
private final SyncEventPayloadReader payloadReader;
@Override
@@ -72,10 +80,10 @@ public void handle(SyncOutboxEvent event) {
Optional latestJob = embeddingJobRepository
.findTopByDocumentVersionIdAndEmbeddingModelIdOrderByIdDesc(version.getId(), model.getId());
if (event.getEventType() == SyncEventType.DOCUMENT_REINDEX_REQUESTED && latestJob.isPresent()) {
- retryOrReuse(latestJob.get());
+ retryOrReuse(latestJob.get(), event, version, model);
return;
}
- if (latestJob.isPresent() || version.getStatus() != DocumentVersionStatus.UPLOADED) {
+ if (latestJob.isPresent() || !RECOVERABLE_PROCESSING_STATUSES.contains(version.getStatus())) {
throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT);
}
@@ -91,7 +99,12 @@ public void handle(SyncOutboxEvent event) {
);
}
- private void retryOrReuse(EmbeddingJob job) {
+ private void retryOrReuse(
+ EmbeddingJob job,
+ SyncOutboxEvent event,
+ DocumentVersion version,
+ EmbeddingModel model
+ ) {
if (job.getStatus() == EmbeddingJobStatus.FAILED) {
embeddingJobManualRetryService.retry(job.getId());
return;
@@ -100,6 +113,10 @@ private void retryOrReuse(EmbeddingJob job) {
|| job.getStatus() == EmbeddingJobStatus.PROCESSING) {
return;
}
+ if (job.getStatus() == EmbeddingJobStatus.INDEXED) {
+ indexedVersionVectorRepairService.repair(version.getId(), model, event.getEventId());
+ return;
+ }
throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT);
}
diff --git a/backend/src/main/resources/application.yml b/backend/src/main/resources/application.yml
index 1118713..e23517e 100644
--- a/backend/src/main/resources/application.yml
+++ b/backend/src/main/resources/application.yml
@@ -67,6 +67,12 @@ sync:
lease-recovery-batch-size: ${SYNC_DISPATCHER_LEASE_RECOVERY_BATCH_SIZE:100}
retry-initial-delay: ${SYNC_DISPATCHER_RETRY_INITIAL_DELAY:5s}
retry-max-delay: ${SYNC_DISPATCHER_RETRY_MAX_DELAY:1m}
+ reconciliation:
+ enabled: ${SYNC_RECONCILIATION_ENABLED:false}
+ mode: ${SYNC_RECONCILIATION_MODE:DRY_RUN}
+ batch-size: ${SYNC_RECONCILIATION_BATCH_SIZE:100}
+ interval: ${SYNC_RECONCILIATION_INTERVAL:5m}
+ stalled-threshold: ${SYNC_RECONCILIATION_STALLED_THRESHOLD:15m}
server:
port: 8080
diff --git a/backend/src/main/resources/db/migration/V37__create_sync_consistency_issues.sql b/backend/src/main/resources/db/migration/V37__create_sync_consistency_issues.sql
new file mode 100644
index 0000000..73e4a04
--- /dev/null
+++ b/backend/src/main/resources/db/migration/V37__create_sync_consistency_issues.sql
@@ -0,0 +1,55 @@
+-- sync_consistency_issues: Reconciler가 발견한 동일 불일치를 하나의 생명주기로 관리한다.
+CREATE TABLE sync_consistency_issues (
+ id BIGSERIAL PRIMARY KEY,
+ issue_key VARCHAR(255) NOT NULL,
+ issue_type VARCHAR(50) NOT NULL,
+ severity VARCHAR(20) NOT NULL,
+ status VARCHAR(20) NOT NULL,
+ document_id BIGINT REFERENCES documents (id),
+ document_version_id BIGINT REFERENCES document_versions (id),
+ embedding_model_id BIGINT REFERENCES embedding_models (id),
+ expected_json TEXT,
+ actual_json TEXT,
+ detected_at TIMESTAMP NOT NULL,
+ last_detected_at TIMESTAMP NOT NULL,
+ repair_event_id UUID REFERENCES sync_outbox_events (event_id),
+ repair_attempt_count INT NOT NULL DEFAULT 0,
+ resolved_at TIMESTAMP,
+ resolution_message TEXT,
+ created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
+
+ CONSTRAINT uk_sync_consistency_issues_issue_key UNIQUE (issue_key)
+);
+
+CREATE INDEX idx_sync_consistency_issues_status_last_detected
+ ON sync_consistency_issues (status, last_detected_at DESC, id DESC);
+CREATE INDEX idx_sync_consistency_issues_document_version
+ ON sync_consistency_issues (document_version_id, status);
+CREATE INDEX idx_sync_consistency_issues_repair_event
+ ON sync_consistency_issues (repair_event_id);
+
+-- sync_reconciliation_runs: Batch 실행의 범위와 탐지·복구 결과를 운영 지표로 보존한다.
+CREATE TABLE sync_reconciliation_runs (
+ id BIGSERIAL PRIMARY KEY,
+ run_id UUID NOT NULL,
+ mode VARCHAR(20) NOT NULL,
+ status VARCHAR(20) NOT NULL,
+ start_cursor BIGINT NOT NULL,
+ end_cursor BIGINT NOT NULL,
+ scanned_count INT NOT NULL,
+ detected_count INT NOT NULL,
+ repair_requested_count INT NOT NULL,
+ started_at TIMESTAMP NOT NULL,
+ completed_at TIMESTAMP,
+ error_code VARCHAR(100),
+ created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
+ updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
+
+ CONSTRAINT uk_sync_reconciliation_runs_run_id UNIQUE (run_id)
+);
+
+CREATE INDEX idx_sync_reconciliation_runs_started_at
+ ON sync_reconciliation_runs (started_at DESC, id DESC);
+CREATE INDEX idx_sync_reconciliation_runs_status
+ ON sync_reconciliation_runs (status, started_at DESC);
diff --git a/backend/src/test/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionServiceTest.java b/backend/src/test/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionServiceTest.java
index f8bc1c8..2f8212d 100644
--- a/backend/src/test/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionServiceTest.java
+++ b/backend/src/test/java/com/opensource/docgrid/domain/embedding/service/command/DocumentEmbeddingTransactionServiceTest.java
@@ -161,17 +161,17 @@ void prepare_replaysCompletedEmbeddings() {
}
@Test
- @DisplayName("일부 Chunk의 Embedding만 저장된 상태는 내부 데이터 모순으로 거부한다")
- void prepare_rejectsPartialEmbeddings() {
+ @DisplayName("Reconciler가 확인한 일부 Embedding Set은 누락 Vector 생성을 재개한다")
+ void prepare_resumesPartialEmbeddings() {
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));
+ PreparationResult result = service.prepare(JOB_ID, ATTEMPT_ID, WORKER_ID, CLAIM_TOKEN);
+
+ assertThat(result.isReplay()).isFalse();
+ assertThat(result.work().chunks()).hasSize(2);
}
@Test
@@ -254,6 +254,36 @@ void complete_savesEmbeddingSet() {
then(indexingEventRepository).shouldHaveNoInteractions();
}
+ @Test
+ @DisplayName("부분 Embedding Set 복구는 기존 Vector를 보존하고 누락 Chunk만 저장한다")
+ @SuppressWarnings("unchecked")
+ void complete_preservesExistingEmbeddingAndSavesMissingChunk() {
+ prepareEntities(DocumentVersionStatus.EMBEDDING);
+ List chunks = List.of(chunk(0, "첫 번째"), chunk(1, "두 번째"));
+ givenValidContext(chunks);
+ given(embeddingRepository.countByDocumentVersionIdAndEmbeddingModelId(VERSION_ID, MODEL_ID))
+ .willReturn(1L);
+ given(embeddingRepository.findChunkIdsByDocumentVersionIdAndEmbeddingModelId(VERSION_ID, MODEL_ID))
+ .willReturn(List.of(chunks.get(0).getId()));
+
+ 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.embeddingCount()).isEqualTo(2);
+ ArgumentCaptor> embeddingsCaptor = ArgumentCaptor.forClass(List.class);
+ then(embeddingRepository).should().saveAllAndFlush(embeddingsCaptor.capture());
+ assertThat(embeddingsCaptor.getValue())
+ .singleElement()
+ .extracting(Embedding::getChunk)
+ .isSameAs(chunks.get(1));
+ }
+
@Test
@DisplayName("완료 단계에서 선행 요청의 전체 저장을 발견하면 Insert 없이 재생한다")
void complete_replaysConcurrentWinner() {
diff --git a/backend/src/test/java/com/opensource/docgrid/domain/embedding/service/command/IndexedVersionVectorRepairServiceTest.java b/backend/src/test/java/com/opensource/docgrid/domain/embedding/service/command/IndexedVersionVectorRepairServiceTest.java
new file mode 100644
index 0000000..90d7e8b
--- /dev/null
+++ b/backend/src/test/java/com/opensource/docgrid/domain/embedding/service/command/IndexedVersionVectorRepairServiceTest.java
@@ -0,0 +1,113 @@
+package com.opensource.docgrid.domain.embedding.service.command;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.BDDMockito.given;
+import static org.mockito.BDDMockito.then;
+import static org.mockito.Mockito.never;
+
+import java.util.Optional;
+import java.util.UUID;
+
+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.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.DocumentRepository;
+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.enums.EmbeddingJobStatus;
+import com.opensource.docgrid.domain.embedding.enums.EmbeddingStatus;
+import com.opensource.docgrid.domain.embedding.fixture.EmbeddingModelFixture;
+import com.opensource.docgrid.domain.embedding.repository.EmbeddingJobRepository;
+import com.opensource.docgrid.domain.embedding.repository.EmbeddingRepository;
+
+/**
+ * INDEXED Vector 복구가 기존 Embedding을 삭제하지 않고 Version 상태와 새 Job만 원자 전환하는지 검증한다.
+ */
+@ExtendWith(MockitoExtension.class)
+@DisplayName("IndexedVersionVectorRepairService 단위 테스트")
+class IndexedVersionVectorRepairServiceTest {
+
+ @Mock private DocumentVersionRepository documentVersionRepository;
+ @Mock private DocumentRepository documentRepository;
+ @Mock private DocumentChunkRepository documentChunkRepository;
+ @Mock private EmbeddingJobRepository embeddingJobRepository;
+ @Mock private EmbeddingRepository embeddingRepository;
+
+ private IndexedVersionVectorRepairService service;
+ private Document document;
+ private DocumentVersion version;
+ private EmbeddingModel model;
+
+ @BeforeEach
+ void setUp() {
+ service = new IndexedVersionVectorRepairService(
+ documentVersionRepository,
+ documentRepository,
+ documentChunkRepository,
+ embeddingJobRepository,
+ embeddingRepository
+ );
+ document = Document.builder()
+ .title("Vector 복구 문서")
+ .documentType(DocumentType.PDF)
+ .sourceType(DocumentSourceType.UPLOAD)
+ .status(DocumentStatus.INDEXED)
+ .visibility(VisibilityType.PRIVATE)
+ .build();
+ ReflectionTestUtils.setField(document, "id", 3L);
+ version = DocumentVersion.builder()
+ .document(document)
+ .versionNo(1)
+ .status(DocumentVersionStatus.INDEXED)
+ .build();
+ ReflectionTestUtils.setField(version, "id", 11L);
+ document.updateCurrentVersion(version);
+ model = EmbeddingModelFixture.createDefaultModel();
+ ReflectionTestUtils.setField(model, "id", 7L);
+
+ given(documentVersionRepository.findByIdForUpdate(11L)).willReturn(Optional.of(version));
+ given(documentRepository.findByIdForUpdate(3L)).willReturn(Optional.of(document));
+ given(documentChunkRepository.existsByDocumentVersionId(11L)).willReturn(true);
+ given(documentChunkRepository.countByDocumentVersionId(11L)).willReturn(3L);
+ given(embeddingJobRepository.existsByDocumentVersionIdAndStatusIn(any(), any())).willReturn(false);
+ given(embeddingRepository.countByDocumentVersionIdAndEmbeddingModelId(11L, 7L)).willReturn(1L);
+ given(embeddingRepository.countByDocumentVersionIdAndEmbeddingModelIdAndStatus(
+ 11L, 7L, EmbeddingStatus.ACTIVE
+ )).willReturn(1L);
+ given(embeddingJobRepository.save(any(EmbeddingJob.class)))
+ .willAnswer(invocation -> invocation.getArgument(0));
+ }
+
+ @Test
+ @DisplayName("부분 누락은 기존 Vector를 보존한 채 CHUNKED와 PENDING Job으로 전환한다")
+ void repair_preservesVectorsAndCreatesJob() {
+ UUID sourceEventId = UUID.randomUUID();
+
+ EmbeddingJob result = service.repair(11L, model, sourceEventId);
+
+ assertThat(version.getStatus()).isEqualTo(DocumentVersionStatus.CHUNKED);
+ assertThat(document.getStatus()).isEqualTo(DocumentStatus.INDEXING);
+ assertThat(result.getStatus()).isEqualTo(EmbeddingJobStatus.PENDING);
+ assertThat(result.getSourceEventId()).isEqualTo(sourceEventId);
+ then(embeddingRepository).should(never()).deleteByDocumentVersionId(any());
+ ArgumentCaptor jobCaptor = ArgumentCaptor.forClass(EmbeddingJob.class);
+ then(embeddingJobRepository).should().save(jobCaptor.capture());
+ assertThat(jobCaptor.getValue().getDocumentVersion()).isSameAs(version);
+ }
+}
diff --git a/backend/src/test/java/com/opensource/docgrid/domain/sync/service/SyncConsistencyInspectorTest.java b/backend/src/test/java/com/opensource/docgrid/domain/sync/service/SyncConsistencyInspectorTest.java
new file mode 100644
index 0000000..6af44af
--- /dev/null
+++ b/backend/src/test/java/com/opensource/docgrid/domain/sync/service/SyncConsistencyInspectorTest.java
@@ -0,0 +1,148 @@
+package com.opensource.docgrid.domain.sync.service;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.BDDMockito.given;
+
+import java.time.Duration;
+import java.time.LocalDateTime;
+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.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.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.embedding.entity.EmbeddingModel;
+import com.opensource.docgrid.domain.embedding.enums.EmbeddingStatus;
+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.sync.config.SyncReconciliationProperties;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencyIssueType;
+
+/**
+ * 정합성 검사가 현재 Version만 Vector 대상으로 삼고 안전한 부분 누락만 복구 가능으로 분류하는지 검증한다.
+ */
+@ExtendWith(MockitoExtension.class)
+@DisplayName("SyncConsistencyInspector 단위 테스트")
+class SyncConsistencyInspectorTest {
+
+ private static final LocalDateTime INSPECTED_AT = LocalDateTime.of(2026, 8, 13, 20, 0);
+
+ @Mock private EmbeddingJobRepository embeddingJobRepository;
+ @Mock private DocumentChunkRepository documentChunkRepository;
+ @Mock private EmbeddingRepository embeddingRepository;
+
+ private SyncConsistencyInspector inspector;
+ private EmbeddingModel model;
+
+ @BeforeEach
+ void setUp() {
+ SyncReconciliationProperties properties = new SyncReconciliationProperties();
+ properties.setStalledThreshold(Duration.ofMinutes(15));
+ inspector = new SyncConsistencyInspector(
+ embeddingJobRepository,
+ documentChunkRepository,
+ embeddingRepository,
+ properties
+ );
+ model = EmbeddingModelFixture.createDefaultModel();
+ ReflectionTestUtils.setField(model, "id", 7L);
+ }
+
+ @Test
+ @DisplayName("현재 Version의 기존 ACTIVE Vector는 보존하고 실제 누락 행만 자동복구 대상으로 분류한다")
+ void inspect_marksPartialMissingEmbeddingsRepairable() {
+ DocumentVersion version = indexedCurrentVersion(11L);
+ given(documentChunkRepository.countByDocumentVersionId(11L)).willReturn(3L);
+ given(embeddingRepository.countByDocumentVersionIdAndEmbeddingModelIdAndStatus(
+ 11L, 7L, EmbeddingStatus.ACTIVE
+ )).willReturn(1L);
+ given(embeddingRepository.countByDocumentVersionIdAndEmbeddingModelId(11L, 7L)).willReturn(1L);
+ given(embeddingRepository.countByDocumentVersionId(11L)).willReturn(1L);
+
+ List observations = inspector.inspect(version, model, INSPECTED_AT);
+
+ assertThat(observations)
+ .singleElement()
+ .satisfies(observation -> {
+ assertThat(observation.issueType()).isEqualTo(SyncConsistencyIssueType.MISSING_EMBEDDINGS);
+ assertThat(observation.repairable()).isTrue();
+ });
+ }
+
+ @Test
+ @DisplayName("행은 모두 있으나 ACTIVE 상태가 깨진 Vector Set은 보고만 한다")
+ void inspect_reportsStatusDamageWithoutAutomaticRepair() {
+ DocumentVersion version = indexedCurrentVersion(12L);
+ given(documentChunkRepository.countByDocumentVersionId(12L)).willReturn(3L);
+ given(embeddingRepository.countByDocumentVersionIdAndEmbeddingModelIdAndStatus(
+ 12L, 7L, EmbeddingStatus.ACTIVE
+ )).willReturn(1L);
+ given(embeddingRepository.countByDocumentVersionIdAndEmbeddingModelId(12L, 7L)).willReturn(3L);
+ given(embeddingRepository.countByDocumentVersionId(12L)).willReturn(3L);
+
+ List observations = inspector.inspect(version, model, INSPECTED_AT);
+
+ assertThat(observations)
+ .singleElement()
+ .satisfies(observation -> assertThat(observation.repairable()).isFalse());
+ }
+
+ @Test
+ @DisplayName("과거 INDEXED Version의 STALE Vector는 현재 검색 상태 불일치로 오인하지 않는다")
+ void inspect_ignoresHistoricalIndexedVersionVectors() {
+ Document document = document(3L, DocumentStatus.INDEXED);
+ DocumentVersion historical = version(document, 13L, DocumentVersionStatus.INDEXED);
+ DocumentVersion current = version(document, 14L, DocumentVersionStatus.INDEXED);
+ document.updateCurrentVersion(current);
+ given(documentChunkRepository.countByDocumentVersionId(13L)).willReturn(3L);
+
+ assertThat(inspector.inspect(historical, model, INSPECTED_AT)).isEmpty();
+ }
+
+ private DocumentVersion indexedCurrentVersion(Long versionId) {
+ Document document = document(versionId, DocumentStatus.INDEXED);
+ DocumentVersion version = version(document, versionId, DocumentVersionStatus.INDEXED);
+ document.updateCurrentVersion(version);
+ return version;
+ }
+
+ private Document document(Long id, DocumentStatus status) {
+ Document document = Document.builder()
+ .title("정합성 검사 문서")
+ .documentType(DocumentType.PDF)
+ .sourceType(DocumentSourceType.UPLOAD)
+ .status(status)
+ .visibility(VisibilityType.PRIVATE)
+ .build();
+ ReflectionTestUtils.setField(document, "id", id);
+ return document;
+ }
+
+ private DocumentVersion version(
+ Document document,
+ Long id,
+ DocumentVersionStatus status
+ ) {
+ DocumentVersion version = DocumentVersion.builder()
+ .document(document)
+ .versionNo(1)
+ .status(status)
+ .build();
+ ReflectionTestUtils.setField(version, "id", id);
+ ReflectionTestUtils.setField(version, "updatedAt", INSPECTED_AT.minusMinutes(1));
+ return version;
+ }
+}
diff --git a/backend/src/test/java/com/opensource/docgrid/domain/sync/service/command/SyncReconciliationBatchServiceTest.java b/backend/src/test/java/com/opensource/docgrid/domain/sync/service/command/SyncReconciliationBatchServiceTest.java
new file mode 100644
index 0000000..d8dd890
--- /dev/null
+++ b/backend/src/test/java/com/opensource/docgrid/domain/sync/service/command/SyncReconciliationBatchServiceTest.java
@@ -0,0 +1,163 @@
+package com.opensource.docgrid.domain.sync.service.command;
+
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.BDDMockito.given;
+import static org.mockito.BDDMockito.then;
+import static org.mockito.Mockito.never;
+
+import java.time.Clock;
+import java.time.Instant;
+import java.time.LocalDateTime;
+import java.time.ZoneId;
+import java.util.List;
+import java.util.UUID;
+
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.DisplayName;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.ExtendWith;
+import org.mockito.Mock;
+import org.mockito.junit.jupiter.MockitoExtension;
+import org.springframework.data.domain.Pageable;
+import org.springframework.test.util.ReflectionTestUtils;
+
+import com.opensource.docgrid.domain.document.entity.DocumentVersion;
+import com.opensource.docgrid.domain.document.enums.DocumentVersionStatus;
+import com.opensource.docgrid.domain.document.repository.DocumentVersionRepository;
+import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel;
+import com.opensource.docgrid.domain.embedding.fixture.EmbeddingModelFixture;
+import com.opensource.docgrid.domain.embedding.service.query.EmbeddingModelQueryService;
+import com.opensource.docgrid.domain.sync.config.SyncReconciliationProperties;
+import com.opensource.docgrid.domain.sync.dto.SyncReconciliationBatchResult;
+import com.opensource.docgrid.domain.sync.entity.SyncConsistencyIssue;
+import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent;
+import com.opensource.docgrid.domain.sync.enums.SyncAggregateType;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencyIssueStatus;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencyIssueType;
+import com.opensource.docgrid.domain.sync.enums.SyncConsistencySeverity;
+import com.opensource.docgrid.domain.sync.enums.SyncEventType;
+import com.opensource.docgrid.domain.sync.enums.SyncReconciliationMode;
+import com.opensource.docgrid.domain.sync.service.SyncConsistencyInspector;
+import com.opensource.docgrid.domain.sync.service.SyncConsistencyObservation;
+import com.opensource.docgrid.domain.sync.service.SyncOrphanInspector;
+
+/**
+ * Cursor Batch가 반복 Issue를 수렴시키고 REPAIR 모드에서만 한 Outbox 복구 요청을 만드는지 검증한다.
+ */
+@ExtendWith(MockitoExtension.class)
+@DisplayName("SyncReconciliationBatchService 단위 테스트")
+class SyncReconciliationBatchServiceTest {
+
+ private static final LocalDateTime NOW = LocalDateTime.of(2026, 8, 13, 21, 0);
+
+ @Mock private DocumentVersionRepository documentVersionRepository;
+ @Mock private EmbeddingModelQueryService embeddingModelQueryService;
+ @Mock private SyncConsistencyInspector syncConsistencyInspector;
+ @Mock private SyncConsistencyIssueService syncConsistencyIssueService;
+ @Mock private SyncOrphanInspector syncOrphanInspector;
+ @Mock private SyncEventWriter syncEventWriter;
+
+ private SyncReconciliationBatchService service;
+ private DocumentVersion version;
+ private EmbeddingModel model;
+ private SyncConsistencyObservation observation;
+ private SyncConsistencyIssue issue;
+
+ @BeforeEach
+ void setUp() {
+ SyncReconciliationProperties properties = new SyncReconciliationProperties();
+ properties.setBatchSize(100);
+ service = new SyncReconciliationBatchService(
+ documentVersionRepository,
+ embeddingModelQueryService,
+ syncConsistencyInspector,
+ syncConsistencyIssueService,
+ syncOrphanInspector,
+ syncEventWriter,
+ properties,
+ Clock.fixed(Instant.parse("2026-08-13T12:00:00Z"), ZoneId.of("Asia/Seoul"))
+ );
+ version = DocumentVersion.builder()
+ .versionNo(1)
+ .status(DocumentVersionStatus.UPLOADED)
+ .build();
+ ReflectionTestUtils.setField(version, "id", 11L);
+ model = EmbeddingModelFixture.createDefaultModel();
+ ReflectionTestUtils.setField(model, "id", 7L);
+ observation = new SyncConsistencyObservation(
+ "MISSING_JOB:VERSION:11:MODEL:7",
+ SyncConsistencyIssueType.MISSING_JOB,
+ SyncConsistencySeverity.ERROR,
+ null,
+ version,
+ model,
+ "{\"liveJob\":true}",
+ "{\"liveJob\":false}",
+ true
+ );
+ issue = SyncConsistencyIssue.builder()
+ .issueKey(observation.issueKey())
+ .issueType(observation.issueType())
+ .severity(observation.severity())
+ .documentVersion(version)
+ .embeddingModel(model)
+ .expectedJson(observation.expectedJson())
+ .actualJson(observation.actualJson())
+ .detectedAt(NOW)
+ .build();
+ given(embeddingModelQueryService.getActiveModel()).willReturn(model);
+ given(documentVersionRepository.findReconciliationBatchAfterId(any(Long.class), any(Pageable.class)))
+ .willReturn(List.of(version));
+ given(syncConsistencyInspector.inspect(version, model, NOW)).willReturn(List.of(observation));
+ given(syncConsistencyIssueService.detect(observation, NOW)).willReturn(issue);
+ }
+
+ @Test
+ @DisplayName("REPAIR 모드는 안전한 OPEN Issue에 Outbox Event 하나를 연결한다")
+ void reconcile_requestsSingleRepairEvent() {
+ UUID runId = UUID.randomUUID();
+ SyncOutboxEvent event = repairEvent();
+ given(syncOrphanInspector.inspect()).willReturn(List.of());
+ given(syncEventWriter.recordDocumentReindexRequested(any(), any(), any()))
+ .willReturn(event);
+
+ SyncReconciliationBatchResult result = service.reconcile(runId, 0L, SyncReconciliationMode.REPAIR);
+
+ assertThat(result.scannedCount()).isEqualTo(1);
+ assertThat(result.detectedCount()).isEqualTo(1);
+ assertThat(result.repairRequestedCount()).isEqualTo(1);
+ assertThat(issue.getStatus()).isEqualTo(SyncConsistencyIssueStatus.REPAIRING);
+ assertThat(issue.getRepairEventId()).isEqualTo(event.getEventId());
+ }
+
+ @Test
+ @DisplayName("DRY_RUN은 Issue만 기록하고 복구 Event를 만들지 않는다")
+ void reconcile_dryRunDoesNotCreateRepairEvent() {
+ SyncReconciliationBatchResult result = service.reconcile(
+ UUID.randomUUID(),
+ 1L,
+ SyncReconciliationMode.DRY_RUN
+ );
+
+ assertThat(result.detectedCount()).isEqualTo(1);
+ assertThat(result.repairRequestedCount()).isZero();
+ assertThat(issue.getStatus()).isEqualTo(SyncConsistencyIssueStatus.OPEN);
+ then(syncEventWriter).should(never()).recordDocumentReindexRequested(any(), any(), any());
+ }
+
+ private SyncOutboxEvent repairEvent() {
+ return SyncOutboxEvent.builder()
+ .eventId(UUID.randomUUID())
+ .idempotencyKey("repair:11")
+ .aggregateType(SyncAggregateType.DOCUMENT_VERSION)
+ .aggregateId(11L)
+ .aggregateVersion(1L)
+ .eventType(SyncEventType.DOCUMENT_REINDEX_REQUESTED)
+ .payloadJson("{\"embeddingModelId\":7}")
+ .availableAt(NOW)
+ .occurredAt(NOW)
+ .maxRetryCount(5)
+ .build();
+ }
+}
diff --git a/backend/src/test/java/com/opensource/docgrid/domain/sync/service/handler/DocumentVersionSyncEventHandlerTest.java b/backend/src/test/java/com/opensource/docgrid/domain/sync/service/handler/DocumentVersionSyncEventHandlerTest.java
index f6f36ce..d988a78 100644
--- a/backend/src/test/java/com/opensource/docgrid/domain/sync/service/handler/DocumentVersionSyncEventHandlerTest.java
+++ b/backend/src/test/java/com/opensource/docgrid/domain/sync/service/handler/DocumentVersionSyncEventHandlerTest.java
@@ -30,6 +30,7 @@
import com.opensource.docgrid.domain.embedding.repository.EmbeddingJobRepository;
import com.opensource.docgrid.domain.embedding.repository.EmbeddingModelRepository;
import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobManualRetryService;
+import com.opensource.docgrid.domain.embedding.service.command.IndexedVersionVectorRepairService;
import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent;
import com.opensource.docgrid.domain.sync.enums.SyncAggregateType;
import com.opensource.docgrid.domain.sync.enums.SyncEventType;
@@ -46,6 +47,7 @@ class DocumentVersionSyncEventHandlerTest {
@Mock private EmbeddingModelRepository embeddingModelRepository;
@Mock private EmbeddingJobRepository embeddingJobRepository;
@Mock private EmbeddingJobManualRetryService embeddingJobManualRetryService;
+ @Mock private IndexedVersionVectorRepairService indexedVersionVectorRepairService;
private DocumentVersionSyncEventHandler handler;
private DocumentVersion version;
@@ -59,6 +61,7 @@ void setUp() {
embeddingModelRepository,
embeddingJobRepository,
embeddingJobManualRetryService,
+ indexedVersionVectorRepairService,
new SyncEventPayloadReader(new ObjectMapper())
);
version = DocumentVersion.builder().versionNo(1).status(DocumentVersionStatus.UPLOADED).build();