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();