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 aa5ea30..1b6d06e 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 @@ -4,6 +4,7 @@ import java.util.Collection; import java.util.List; import java.util.Optional; +import java.util.UUID; import jakarta.persistence.LockModeType; @@ -25,6 +26,13 @@ */ public interface EmbeddingJobRepository extends JpaRepository { + Optional findBySourceEventId(UUID sourceEventId); + + Optional findTopByDocumentVersionIdAndEmbeddingModelIdOrderByIdDesc( + Long documentVersionId, + Long embeddingModelId + ); + /** * 관리자 목록 화면에 필요한 연관관계를 함께 조회하면서 선택 필터와 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 dd3727f..971024a 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 @@ -48,6 +48,19 @@ int markActiveAsStaleByDocumentVersionId( @Param("documentVersionId") Long documentVersionId ); + /** + * 삭제 문서의 모든 ACTIVE Vector를 즉시 검색 대상에서 제외한다. + */ + @Modifying(flushAutomatically = true) + @Query(value = """ + UPDATE embeddings + SET status = 'STALE', + updated_at = CURRENT_TIMESTAMP + WHERE document_id = :documentId + AND status = 'ACTIVE' + """, nativeQuery = true) + int markActiveAsStaleByDocumentId(@Param("documentId") Long documentId); + /** * 수동 재처리로 다시 인덱싱할 Version의 Embedding 행을 한 SQL로 제거한다. * diff --git a/backend/src/main/java/com/opensource/docgrid/domain/permission/service/command/CollectionPermissionCommandService.java b/backend/src/main/java/com/opensource/docgrid/domain/permission/service/command/CollectionPermissionCommandService.java index 8421263..8d8d476 100644 --- a/backend/src/main/java/com/opensource/docgrid/domain/permission/service/command/CollectionPermissionCommandService.java +++ b/backend/src/main/java/com/opensource/docgrid/domain/permission/service/command/CollectionPermissionCommandService.java @@ -27,6 +27,8 @@ import com.opensource.docgrid.domain.user.repository.DepartmentRepository; import com.opensource.docgrid.domain.user.repository.RoleRepository; import com.opensource.docgrid.domain.user.repository.UserRepository; +import com.opensource.docgrid.domain.sync.enums.SyncPermissionOperation; +import com.opensource.docgrid.domain.sync.service.command.SyncEventWriter; import com.opensource.docgrid.global.exception.DocGridException; import com.opensource.docgrid.global.exception.ErrorCode; @@ -46,6 +48,7 @@ public class CollectionPermissionCommandService { private final DepartmentRepository departmentRepository; private final PermissionConverter permissionConverter; private final PermissionQueryService permissionQueryService; + private final SyncEventWriter syncEventWriter; // 컬렉션 권한 부여 public CollectionPermissionResponse grantPermission(Long collectionId, Long grantorId, @@ -99,6 +102,13 @@ public CollectionPermissionResponse grantPermission(Long collectionId, Long gran updateCacheForCollection(collectionId, targetUser, permissions, permission.getId(), request.expiresAt()); } + // 권한 원장과 같은 Transaction에 컬렉션 캐시 재투영 의도를 기록한다. + syncEventWriter.recordPermissionCacheRefresh( + AccessSourceType.DIRECT_COLLECTION_PERMISSION, + permission.getId(), + SyncPermissionOperation.GRANTED + ); + return permissionConverter.toCollectionPermissionResponse(permission); } @@ -121,6 +131,11 @@ public void revokePermission(Long collectionId, Long permissionId, Long revokerI } collectionPermissionRepository.delete(permission); + syncEventWriter.recordPermissionCacheRefresh( + AccessSourceType.DIRECT_COLLECTION_PERMISSION, + permissionId, + SyncPermissionOperation.REVOKED + ); } // targetType과 ID 필드 조합 유효성 검사 diff --git a/backend/src/main/java/com/opensource/docgrid/domain/permission/service/command/DocumentPermissionCommandService.java b/backend/src/main/java/com/opensource/docgrid/domain/permission/service/command/DocumentPermissionCommandService.java index efc3f77..2511309 100644 --- a/backend/src/main/java/com/opensource/docgrid/domain/permission/service/command/DocumentPermissionCommandService.java +++ b/backend/src/main/java/com/opensource/docgrid/domain/permission/service/command/DocumentPermissionCommandService.java @@ -22,6 +22,8 @@ import com.opensource.docgrid.domain.user.repository.DepartmentRepository; import com.opensource.docgrid.domain.user.repository.RoleRepository; import com.opensource.docgrid.domain.user.repository.UserRepository; +import com.opensource.docgrid.domain.sync.enums.SyncPermissionOperation; +import com.opensource.docgrid.domain.sync.service.command.SyncEventWriter; import com.opensource.docgrid.global.exception.DocGridException; import com.opensource.docgrid.global.exception.ErrorCode; @@ -40,6 +42,7 @@ public class DocumentPermissionCommandService { private final DepartmentRepository departmentRepository; private final PermissionConverter permissionConverter; private final PermissionQueryService permissionQueryService; + private final SyncEventWriter syncEventWriter; // 문서 단건 예외 권한 부여 public DocumentPermissionResponse grantPermission(Long documentId, Long grantorId, @@ -94,6 +97,13 @@ public DocumentPermissionResponse grantPermission(Long documentId, Long grantorI AccessSourceType.DIRECT_DOCUMENT_PERMISSION, permission.getId(), request.expiresAt()); } + // 권한 원장과 같은 Transaction에 캐시 재투영 의도를 남겨 후속 누락을 복구할 수 있게 한다. + syncEventWriter.recordPermissionCacheRefresh( + AccessSourceType.DIRECT_DOCUMENT_PERMISSION, + permission.getId(), + SyncPermissionOperation.GRANTED + ); + return permissionConverter.toDocumentPermissionResponse(permission); } @@ -116,6 +126,11 @@ public void revokePermission(Long documentId, Long permissionId, Long revokerId) } documentPermissionRepository.delete(permission); + syncEventWriter.recordPermissionCacheRefresh( + AccessSourceType.DIRECT_DOCUMENT_PERMISSION, + permissionId, + SyncPermissionOperation.REVOKED + ); } // targetType과 ID 필드 조합 유효성 검사 diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/config/SyncDispatcherProperties.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/config/SyncDispatcherProperties.java new file mode 100644 index 0000000..99d1525 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/config/SyncDispatcherProperties.java @@ -0,0 +1,69 @@ +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 jakarta.validation.constraints.AssertTrue; +import jakarta.validation.constraints.Min; +import jakarta.validation.constraints.NotBlank; +import jakarta.validation.constraints.NotNull; +import lombok.Getter; +import lombok.Setter; + +/** + * Sync Dispatcher의 실행 여부, Polling, Lease, 복구 Batch와 Retry Backoff 설정을 바인딩한다. + * + *

API 전용 인스턴스에서는 enabled를 false로 유지할 수 있고, 다중 Dispatcher는 서로 다른 name을 + * 사용해 Event의 현재 소유 인스턴스를 운영 화면에서 식별한다. + */ +@Getter +@Setter +@Validated +@Component +@ConfigurationProperties(prefix = "sync.dispatcher") +public class SyncDispatcherProperties { + + private boolean enabled = false; + + @NotBlank + private String name = "sync-dispatcher"; + + @NotNull + private Duration pollingInterval = Duration.ofSeconds(1); + + @NotNull + private Duration leaseDuration = Duration.ofSeconds(30); + + @NotNull + private Duration leaseRecoveryInterval = Duration.ofSeconds(10); + + @Min(1) + private int leaseRecoveryBatchSize = 100; + + @NotNull + private Duration retryInitialDelay = Duration.ofSeconds(5); + + @NotNull + private Duration retryMaxDelay = Duration.ofMinutes(1); + + @AssertTrue(message = "Sync Dispatcher의 Polling과 Lease 시간은 0보다 커야 합니다.") + public boolean isTimingValid() { + return isPositive(pollingInterval) + && isPositive(leaseDuration) + && isPositive(leaseRecoveryInterval); + } + + @AssertTrue(message = "Sync Retry 최대 지연은 양수인 초기 지연보다 짧을 수 없습니다.") + public boolean isRetryDelayValid() { + return isPositive(retryInitialDelay) + && retryMaxDelay != null + && retryMaxDelay.compareTo(retryInitialDelay) >= 0; + } + + 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/SyncSchedulingConfig.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/config/SyncSchedulingConfig.java new file mode 100644 index 0000000..e1741ca --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/config/SyncSchedulingConfig.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; + +/** + * Sync Dispatcher가 활성화된 실행 인스턴스에서만 Polling과 Lease Recovery Scheduler를 켠다. + */ +@Configuration +@EnableScheduling +@ConditionalOnProperty(prefix = "sync.dispatcher", name = "enabled", havingValue = "true") +public class SyncSchedulingConfig { +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/dto/ClaimedSyncEvent.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/dto/ClaimedSyncEvent.java new file mode 100644 index 0000000..16b7e00 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/dto/ClaimedSyncEvent.java @@ -0,0 +1,12 @@ +package com.opensource.docgrid.domain.sync.dto; + +import java.util.UUID; + +/** + * Dispatcher가 후속 처리 Transaction에 전달하는 최소 Event 소유권 Snapshot이다. + * + *

Payload와 도메인 상태는 Dispatch 시 DB에서 다시 읽고, 이 DTO는 Event ID와 현재 Claim Token만 + * 전달해 오래된 Scheduler 실행이 새 소유권으로 작업하지 못하게 한다. + */ +public record ClaimedSyncEvent(UUID eventId, UUID claimToken) { +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/entity/SyncOutboxEvent.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/entity/SyncOutboxEvent.java index bb5fed9..a64cbf6 100644 --- a/backend/src/main/java/com/opensource/docgrid/domain/sync/entity/SyncOutboxEvent.java +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/entity/SyncOutboxEvent.java @@ -141,4 +141,167 @@ public SyncOutboxEvent( this.occurredAt = occurredAt; this.maxRetryCount = maxRetryCount; } + + /** + * 실행 가능한 PENDING Event에 Dispatcher Lease 소유권을 원자적으로 부여한다. + */ + public void claim( + String dispatcherName, + UUID newClaimToken, + LocalDateTime claimedAt, + LocalDateTime newLockExpiresAt + ) { + // 1. 이미 Claim됐거나 끝난 Event의 소유권을 덮어쓰지 않는다. + if (status != SyncEventStatus.PENDING || availableAt == null || availableAt.isAfter(claimedAt)) { + throw new IllegalStateException("실행 가능한 PENDING Event만 Claim할 수 있습니다."); + } + // 2. 빈 소유자나 유효하지 않은 Lease가 영속화되지 않게 입력 계약을 검증한다. + if (dispatcherName == null + || dispatcherName.isBlank() + || newClaimToken == null + || claimedAt == null + || newLockExpiresAt == null + || !newLockExpiresAt.isAfter(claimedAt)) { + throw new IllegalArgumentException("유효한 Dispatcher 소유권과 Lease가 필요합니다."); + } + + // 3. 처리 상태와 현재 Claim 세대의 소유권을 함께 반영한다. + status = SyncEventStatus.PROCESSING; + lockedBy = dispatcherName; + claimToken = newClaimToken; + lockExpiresAt = newLockExpiresAt; + } + + /** + * 현재 Claim이 Handler 부작용까지 Commit할 준비가 됐을 때 Event를 완료한다. + */ + public void complete(UUID currentClaimToken, LocalDateTime completedAt) { + validateActiveOwnership(currentClaimToken, completedAt); + status = SyncEventStatus.PROCESSED; + processedAt = completedAt; + lastErrorCode = null; + lastErrorMessage = null; + clearOwnership(); + } + + /** + * 장시간 Handler가 현재 Claim 세대를 유지한 채 Lease 만료 시각만 연장한다. + */ + public void renewLease( + UUID currentClaimToken, + LocalDateTime renewedAt, + LocalDateTime renewedLockExpiresAt + ) { + validateActiveOwnership(currentClaimToken, renewedAt); + if (renewedLockExpiresAt == null || !renewedLockExpiresAt.isAfter(lockExpiresAt)) { + throw new IllegalArgumentException("새 Lease 만료 시각은 현재 Lease보다 늦어야 합니다."); + } + lockExpiresAt = renewedLockExpiresAt; + } + + /** + * 현재 처리 실패를 기록하고 지정 시각 이후 다시 Claim 가능한 Queue 상태로 되돌린다. + */ + public void scheduleRetry( + UUID currentClaimToken, + String errorCode, + String errorMessage, + LocalDateTime failedAt, + LocalDateTime nextAvailableAt + ) { + validateActiveOwnership(currentClaimToken, failedAt); + if (retryCount + 1 >= maxRetryCount || nextAvailableAt == null || nextAvailableAt.isBefore(failedAt)) { + throw new IllegalStateException("남은 Retry와 다음 실행 시각이 필요합니다."); + } + status = SyncEventStatus.PENDING; + retryCount++; + availableAt = nextAvailableAt; + lastErrorCode = errorCode; + lastErrorMessage = errorMessage; + clearOwnership(); + } + + /** + * Retry를 모두 소진한 현재 Claim을 최종 실패 상태로 종결한다. + */ + public void markFailed( + UUID currentClaimToken, + String errorCode, + String errorMessage, + LocalDateTime failedAt + ) { + validateActiveOwnership(currentClaimToken, failedAt); + status = SyncEventStatus.FAILED; + retryCount++; + lastErrorCode = errorCode; + lastErrorMessage = errorMessage; + clearOwnership(); + } + + /** + * 관리자가 최종 실패 Event에 실행 기회 한 번을 추가해 Queue로 되돌린다. + */ + public void requeueFailed(LocalDateTime requeuedAt) { + if (status != SyncEventStatus.FAILED || requeuedAt == null) { + throw new IllegalStateException("FAILED Event만 수동 재처리할 수 있습니다."); + } + status = SyncEventStatus.PENDING; + availableAt = requeuedAt; + maxRetryCount++; + lastErrorCode = null; + lastErrorMessage = null; + processedAt = null; + clearOwnership(); + } + + /** + * 만료된 PROCESSING Lease를 Retry Queue 또는 최종 실패 상태로 회수한다. + */ + public void recoverExpiredLease( + String errorCode, + String errorMessage, + LocalDateTime recoveredAt, + LocalDateTime nextAvailableAt + ) { + // 1. 아직 유효하거나 이미 다른 흐름이 끝낸 Event를 오래된 복구 Snapshot으로 변경하지 않는다. + if (status != SyncEventStatus.PROCESSING + || claimToken == null + || lockExpiresAt == null + || recoveredAt == null + || lockExpiresAt.isAfter(recoveredAt)) { + throw new IllegalStateException("만료된 PROCESSING Event만 회수할 수 있습니다."); + } + + // 2. 이번 만료를 실패 횟수에 반영하고 남은 기회에 따라 Queue 또는 최종 실패로 전환한다. + retryCount++; + lastErrorCode = errorCode; + lastErrorMessage = errorMessage; + if (retryCount >= maxRetryCount) { + status = SyncEventStatus.FAILED; + } else { + if (nextAvailableAt == null || nextAvailableAt.isBefore(recoveredAt)) { + throw new IllegalArgumentException("다음 실행 시각은 복구 시각보다 빠를 수 없습니다."); + } + status = SyncEventStatus.PENDING; + availableAt = nextAvailableAt; + } + clearOwnership(); + } + + private void validateActiveOwnership(UUID currentClaimToken, LocalDateTime operatedAt) { + if (status != SyncEventStatus.PROCESSING + || currentClaimToken == null + || !currentClaimToken.equals(claimToken) + || operatedAt == null + || lockExpiresAt == null + || !lockExpiresAt.isAfter(operatedAt)) { + throw new IllegalStateException("유효한 현재 Sync Event 소유권이 필요합니다."); + } + } + + private void clearOwnership() { + lockedBy = null; + claimToken = null; + lockExpiresAt = null; + } } diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncPermissionOperation.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncPermissionOperation.java new file mode 100644 index 0000000..089acfd --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncPermissionOperation.java @@ -0,0 +1,9 @@ +package com.opensource.docgrid.domain.sync.enums; + +/** + * 권한 원장 변경 후 접근 캐시가 반영해야 하는 동작을 구분한다. + */ +public enum SyncPermissionOperation { + GRANTED, + REVOKED +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/lifecycle/SyncEventLeaseRecoveryScheduler.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/lifecycle/SyncEventLeaseRecoveryScheduler.java new file mode 100644 index 0000000..7fcd1c6 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/lifecycle/SyncEventLeaseRecoveryScheduler.java @@ -0,0 +1,64 @@ +package com.opensource.docgrid.domain.sync.lifecycle; + +import java.time.Clock; +import java.time.LocalDateTime; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.atomic.AtomicBoolean; + +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.SyncDispatcherProperties; +import com.opensource.docgrid.domain.sync.repository.SyncOutboxEventRepository; +import com.opensource.docgrid.domain.sync.service.command.SyncEventLeaseRecoveryService; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; + +/** + * 만료된 PROCESSING Sync Event를 제한 Batch로 찾아 독립 복구 Transaction에 전달한다. + */ +@Slf4j +@Component +@RequiredArgsConstructor +@ConditionalOnProperty(prefix = "sync.dispatcher", name = "enabled", havingValue = "true") +public class SyncEventLeaseRecoveryScheduler { + + private final SyncOutboxEventRepository syncOutboxEventRepository; + private final SyncEventLeaseRecoveryService syncEventLeaseRecoveryService; + private final SyncDispatcherProperties syncDispatcherProperties; + private final Clock clock; + private final AtomicBoolean recovering = new AtomicBoolean(false); + + @Scheduled( + fixedDelayString = "${sync.dispatcher.lease-recovery-interval:10s}", + initialDelayString = "${sync.dispatcher.lease-recovery-interval:10s}" + ) + public void recoverExpiredLeases() { + if (!recovering.compareAndSet(false, true)) { + return; + } + try { + LocalDateTime recoveredAt = LocalDateTime.now(clock); + List eventIds = syncOutboxEventRepository.findExpiredProcessingEventIds( + recoveredAt, + syncDispatcherProperties.getLeaseRecoveryBatchSize() + ); + for (UUID eventId : eventIds) { + try { + syncEventLeaseRecoveryService.recover(eventId, recoveredAt); + } catch (RuntimeException exception) { + log.error( + "만료 Sync Event Lease 복구에 실패했습니다. eventId={}, errorType={}", + eventId, + exception.getClass().getSimpleName() + ); + } + } + } finally { + recovering.set(false); + } + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/lifecycle/SyncEventPollingScheduler.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/lifecycle/SyncEventPollingScheduler.java new file mode 100644 index 0000000..0b42ba2 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/lifecycle/SyncEventPollingScheduler.java @@ -0,0 +1,86 @@ +package com.opensource.docgrid.domain.sync.lifecycle; + +import java.util.Optional; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.scheduling.annotation.Scheduled; +import org.springframework.stereotype.Component; + +import com.opensource.docgrid.domain.sync.dto.ClaimedSyncEvent; +import com.opensource.docgrid.domain.sync.service.command.SyncEventClaimService; +import com.opensource.docgrid.domain.sync.service.command.SyncEventDispatchService; +import com.opensource.docgrid.domain.sync.service.command.SyncEventFailureService; +import com.opensource.docgrid.global.exception.DocGridException; + +import lombok.RequiredArgsConstructor; +import lombok.extern.slf4j.Slf4j; + +/** + * PENDING Sync Event를 하나씩 Claim해 멱등 Handler로 전달하는 Polling Scheduler다. + * + *

한 인스턴스의 이전 주기가 끝나기 전에 다음 주기가 겹치지 않으며, Handler 실패는 Rollback 후 별도 + * 실패 Transaction으로 Retry 예약을 시도한다. + */ +@Slf4j +@Component +@RequiredArgsConstructor +@ConditionalOnProperty(prefix = "sync.dispatcher", name = "enabled", havingValue = "true") +public class SyncEventPollingScheduler { + + private static final String HANDLER_FAILURE_MESSAGE = "Sync Event Handler 실행에 실패했습니다."; + + private final SyncEventClaimService syncEventClaimService; + private final SyncEventDispatchService syncEventDispatchService; + private final SyncEventFailureService syncEventFailureService; + private final AtomicBoolean polling = new AtomicBoolean(false); + + @Scheduled( + fixedDelayString = "${sync.dispatcher.polling-interval:1s}", + initialDelayString = "${sync.dispatcher.polling-interval:1s}" + ) + public void poll() { + if (!polling.compareAndSet(false, true)) { + return; + } + try { + Optional claimedEvent = syncEventClaimService.claim(); + claimedEvent.ifPresent(this::dispatch); + } catch (RuntimeException exception) { + log.error("Sync Event Claim에 실패했습니다. errorCode={}", diagnosticCode(exception)); + } finally { + polling.set(false); + } + } + + private void dispatch(ClaimedSyncEvent claimedEvent) { + try { + syncEventDispatchService.dispatch(claimedEvent); + } catch (RuntimeException exception) { + String errorCode = diagnosticCode(exception); + log.error("Sync Event 처리에 실패했습니다. eventId={}, errorCode={}", claimedEvent.eventId(), errorCode); + try { + syncEventFailureService.recordFailure( + claimedEvent.eventId(), + claimedEvent.claimToken(), + errorCode, + HANDLER_FAILURE_MESSAGE + ); + } catch (RuntimeException failureTransitionException) { + // Lease가 먼저 만료됐다면 Recovery Scheduler가 회수하므로 현재 주기에서 소유권을 덮지 않는다. + log.error( + "Sync Event 실패 상태 기록에 실패했습니다. eventId={}, errorCode={}", + claimedEvent.eventId(), + diagnosticCode(failureTransitionException) + ); + } + } + } + + 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/repository/SyncOutboxEventRepository.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/repository/SyncOutboxEventRepository.java index 50b98cf..fca3f12 100644 --- a/backend/src/main/java/com/opensource/docgrid/domain/sync/repository/SyncOutboxEventRepository.java +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/repository/SyncOutboxEventRepository.java @@ -1,20 +1,83 @@ package com.opensource.docgrid.domain.sync.repository; import java.util.Optional; +import java.util.List; +import java.time.LocalDateTime; import java.util.UUID; +import jakarta.persistence.LockModeType; + import org.springframework.data.jpa.repository.JpaRepository; +import org.springframework.data.jpa.repository.Lock; +import org.springframework.data.jpa.repository.Query; +import org.springframework.data.repository.query.Param; import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; /** * 동기화 Outbox Event의 영속성과 멱등 식별자 조회를 담당한다. * - *

Dispatcher의 Claim·Lease Query는 해당 동작을 구현하는 단계에서 이 Repository에 추가한다. + *

일반 멱등 조회와 함께 PostgreSQL의 {@code FOR UPDATE SKIP LOCKED}를 이용해 Dispatcher Claim과 + * 만료 Lease 복구 경쟁을 직렬화한다. */ public interface SyncOutboxEventRepository extends JpaRepository { Optional findByEventId(UUID eventId); Optional findByIdempotencyKey(String idempotencyKey); + + /** + * 현재 실행 가능한 가장 오래된 PENDING Event 한 건을 다른 Dispatcher의 잠금을 기다리지 않고 Claim한다. + */ + @Query(value = """ + SELECT event.* + FROM sync_outbox_events event + WHERE event.status = 'PENDING' + AND event.available_at <= :claimedAt + ORDER BY event.available_at ASC, event.occurred_at ASC, event.id ASC + LIMIT 1 + FOR UPDATE SKIP LOCKED + """, nativeQuery = true) + Optional findNextPendingForUpdate(@Param("claimedAt") LocalDateTime claimedAt); + + /** + * Handler 실행·완료 전이와 다른 소유권 변경을 Event 행에서 직렬화한다. + */ + @Lock(LockModeType.PESSIMISTIC_WRITE) + @Query("SELECT event FROM SyncOutboxEvent event WHERE event.eventId = :eventId") + Optional findByEventIdForUpdate(@Param("eventId") UUID eventId); + + /** + * 만료된 PROCESSING Event 식별자를 오래 만료된 순서대로 제한 조회한다. + */ + @Query(value = """ + SELECT event.event_id + FROM sync_outbox_events event + WHERE event.status = 'PROCESSING' + AND event.lock_expires_at IS NOT NULL + AND event.lock_expires_at <= :recoveredAt + ORDER BY event.lock_expires_at ASC, event.id ASC + LIMIT :batchSize + """, nativeQuery = true) + List findExpiredProcessingEventIds( + @Param("recoveredAt") LocalDateTime recoveredAt, + @Param("batchSize") int batchSize + ); + + /** + * 후보 Event가 여전히 만료 상태일 때만 복구 Transaction의 쓰기 잠금을 획득한다. + */ + @Query(value = """ + SELECT event.* + FROM sync_outbox_events event + WHERE event.event_id = :eventId + AND event.status = 'PROCESSING' + AND event.lock_expires_at IS NOT NULL + AND event.lock_expires_at <= :recoveredAt + FOR UPDATE SKIP LOCKED + """, nativeQuery = true) + Optional findExpiredByEventIdForUpdateSkipLocked( + @Param("eventId") UUID eventId, + @Param("recoveredAt") LocalDateTime recoveredAt + ); } diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncEventHandler.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncEventHandler.java new file mode 100644 index 0000000..7db1bb8 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncEventHandler.java @@ -0,0 +1,19 @@ +package com.opensource.docgrid.domain.sync.service; + +import java.util.Set; + +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncEventType; + +/** + * Sync Event Type별 멱등 부작용을 수행하는 Handler 경계다. + * + *

구현체는 현재 Dispatch Transaction에 참여하며 외부 Commit을 직접 만들지 않는다. 같은 Event가 + * 반복 호출돼도 최종 도메인 상태가 한 번 처리한 결과와 같아야 한다. + */ +public interface SyncEventHandler { + + Set supportedTypes(); + + void handle(SyncOutboxEvent event); +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncEventHandlerRegistry.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncEventHandlerRegistry.java new file mode 100644 index 0000000..8ef685d --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncEventHandlerRegistry.java @@ -0,0 +1,42 @@ +package com.opensource.docgrid.domain.sync.service; + +import java.util.EnumMap; +import java.util.List; +import java.util.Map; + +import org.springframework.stereotype.Component; + +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncEventType; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +/** + * Event Type을 유일한 SyncEventHandler에 연결하고 Dispatch 진입점을 제공한다. + * + *

같은 Type을 둘 이상의 Handler가 선언하면 시작 단계에서 실패해 중복 부작용 가능성을 숨기지 않는다. + */ +@Component +public class SyncEventHandlerRegistry { + + private final Map handlers = new EnumMap<>(SyncEventType.class); + + public SyncEventHandlerRegistry(List candidates) { + for (SyncEventHandler candidate : candidates) { + for (SyncEventType eventType : candidate.supportedTypes()) { + SyncEventHandler previous = handlers.putIfAbsent(eventType, candidate); + if (previous != null) { + throw new IllegalStateException("Sync Event Type별 Handler는 하나만 존재해야 합니다: " + eventType); + } + } + } + } + + public void handle(SyncOutboxEvent event) { + SyncEventHandler handler = handlers.get(event.getEventType()); + if (handler == null) { + throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT); + } + handler.handle(event); + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncEventPayloadReader.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncEventPayloadReader.java new file mode 100644 index 0000000..367a075 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncEventPayloadReader.java @@ -0,0 +1,53 @@ +package com.opensource.docgrid.domain.sync.service; + +import org.springframework.stereotype.Component; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +import lombok.RequiredArgsConstructor; + +/** + * Sync Event의 최소 JSON Payload를 검증해 Handler에 타입 안전한 값으로 제공한다. + * + *

Payload는 최신 도메인 상태가 아니라 조회에 필요한 힌트이므로, Handler는 이 값을 사용해 Entity를 + * 다시 조회하고 관계를 검증해야 한다. + */ +@Component +@RequiredArgsConstructor +public class SyncEventPayloadReader { + + private final ObjectMapper objectMapper; + + public Long requiredLong(SyncOutboxEvent event, String fieldName) { + JsonNode node = read(event).get(fieldName); + if (node == null || !node.canConvertToLong() || node.longValue() <= 0) { + throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT); + } + return node.longValue(); + } + + public String requiredText(SyncOutboxEvent event, String fieldName) { + JsonNode node = read(event).get(fieldName); + if (node == null || !node.isTextual() || node.textValue().isBlank()) { + throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT); + } + return node.textValue(); + } + + private JsonNode read(SyncOutboxEvent event) { + try { + JsonNode root = objectMapper.readTree(event.getPayloadJson()); + if (root == null || !root.isObject()) { + throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT); + } + return root; + } catch (JsonProcessingException | IllegalArgumentException exception) { + throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT, exception); + } + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncEventRetrySchedule.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncEventRetrySchedule.java new file mode 100644 index 0000000..f9378ac --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/SyncEventRetrySchedule.java @@ -0,0 +1,36 @@ +package com.opensource.docgrid.domain.sync.service; + +import java.time.Duration; +import java.time.LocalDateTime; + +import org.springframework.stereotype.Component; + +import com.opensource.docgrid.domain.sync.config.SyncDispatcherProperties; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; + +import lombok.RequiredArgsConstructor; + +/** + * Sync Event 실패 횟수에 따른 지수 Backoff를 설정 상한 안에서 계산한다. + */ +@Component +@RequiredArgsConstructor +public class SyncEventRetrySchedule { + + private final SyncDispatcherProperties syncDispatcherProperties; + + public LocalDateTime nextAvailableAt(SyncOutboxEvent event, LocalDateTime failedAt) { + Duration delay = syncDispatcherProperties.getRetryInitialDelay(); + for (int index = 0; index < event.getRetryCount(); index++) { + if (delay.compareTo(syncDispatcherProperties.getRetryMaxDelay()) >= 0) { + delay = syncDispatcherProperties.getRetryMaxDelay(); + break; + } + delay = delay.multipliedBy(2); + } + if (delay.compareTo(syncDispatcherProperties.getRetryMaxDelay()) > 0) { + delay = syncDispatcherProperties.getRetryMaxDelay(); + } + return failedAt.plus(delay); + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventClaimService.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventClaimService.java new file mode 100644 index 0000000..bca9bb0 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventClaimService.java @@ -0,0 +1,49 @@ +package com.opensource.docgrid.domain.sync.service.command; + +import java.time.Clock; +import java.time.LocalDateTime; +import java.util.Optional; +import java.util.UUID; + +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +import com.opensource.docgrid.domain.sync.config.SyncDispatcherProperties; +import com.opensource.docgrid.domain.sync.dto.ClaimedSyncEvent; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.repository.SyncOutboxEventRepository; + +import lombok.RequiredArgsConstructor; + +/** + * 실행 가능한 PENDING Outbox Event 하나를 잠그고 현재 Dispatcher에 Lease 소유권을 부여한다. + * + *

행 선택과 PROCESSING 전이를 한 Transaction에서 수행해 다중 Scheduler가 같은 Event를 Claim하지 + * 못하게 한다. + */ +@Service +@RequiredArgsConstructor +@Transactional +public class SyncEventClaimService { + + private final SyncOutboxEventRepository syncOutboxEventRepository; + private final SyncDispatcherProperties syncDispatcherProperties; + private final Clock clock; + + public Optional claim() { + LocalDateTime claimedAt = LocalDateTime.now(clock); + return syncOutboxEventRepository.findNextPendingForUpdate(claimedAt) + .map(event -> claim(event, claimedAt)); + } + + private ClaimedSyncEvent claim(SyncOutboxEvent event, LocalDateTime claimedAt) { + UUID claimToken = UUID.randomUUID(); + event.claim( + syncDispatcherProperties.getName(), + claimToken, + claimedAt, + claimedAt.plus(syncDispatcherProperties.getLeaseDuration()) + ); + return new ClaimedSyncEvent(event.getEventId(), claimToken); + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventDispatchService.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventDispatchService.java new file mode 100644 index 0000000..1d0d521 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventDispatchService.java @@ -0,0 +1,59 @@ +package com.opensource.docgrid.domain.sync.service.command; + +import java.time.Clock; +import java.time.LocalDateTime; +import java.util.Objects; + +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; + +import com.opensource.docgrid.domain.sync.dto.ClaimedSyncEvent; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncEventStatus; +import com.opensource.docgrid.domain.sync.repository.SyncOutboxEventRepository; +import com.opensource.docgrid.domain.sync.service.SyncEventHandlerRegistry; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +import lombok.RequiredArgsConstructor; + +/** + * 현재 Claim을 재검증한 뒤 Handler 부작용과 Event 완료를 하나의 독립 Transaction에서 확정한다. + * + *

Handler가 실패하면 전체 Transaction이 Rollback되므로 Event는 PROCESSING에 남고, 호출 Scheduler가 + * 별도 실패 전이 Service로 Retry를 예약한다. + */ +@Service +@RequiredArgsConstructor +public class SyncEventDispatchService { + + private final SyncOutboxEventRepository syncOutboxEventRepository; + private final SyncEventHandlerRegistry syncEventHandlerRegistry; + private final Clock clock; + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void dispatch(ClaimedSyncEvent claimedEvent) { + LocalDateTime dispatchedAt = LocalDateTime.now(clock); + SyncOutboxEvent event = syncOutboxEventRepository.findByEventIdForUpdate(claimedEvent.eventId()) + .orElseThrow(() -> new DocGridException(ErrorCode.SYNC_EVENT_NOT_FOUND)); + validateOwnership(event, claimedEvent, dispatchedAt); + + // Handler 부작용과 완료 상태가 같은 Commit 경계를 공유해야 부분 완료가 남지 않는다. + syncEventHandlerRegistry.handle(event); + event.complete(claimedEvent.claimToken(), LocalDateTime.now(clock)); + } + + private void validateOwnership( + SyncOutboxEvent event, + ClaimedSyncEvent claimedEvent, + LocalDateTime dispatchedAt + ) { + if (event.getStatus() != SyncEventStatus.PROCESSING + || !Objects.equals(event.getClaimToken(), claimedEvent.claimToken()) + || event.getLockExpiresAt() == null + || !event.getLockExpiresAt().isAfter(dispatchedAt)) { + throw new DocGridException(ErrorCode.SYNC_EVENT_OWNERSHIP_INVALID); + } + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventFailureService.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventFailureService.java new file mode 100644 index 0000000..3d71478 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventFailureService.java @@ -0,0 +1,67 @@ +package com.opensource.docgrid.domain.sync.service.command; + +import java.time.Clock; +import java.time.LocalDateTime; +import java.util.Objects; +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.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncEventStatus; +import com.opensource.docgrid.domain.sync.repository.SyncOutboxEventRepository; +import com.opensource.docgrid.domain.sync.service.SyncEventRetrySchedule; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +import lombok.RequiredArgsConstructor; + +/** + * Rollback된 Handler 실행의 실패를 별도 Transaction에서 Retry 또는 최종 실패 상태로 기록한다. + * + *

호출자는 민감한 예외 Message 대신 안정적인 오류 코드와 제한된 진단 문구만 전달해야 한다. + */ +@Service +@RequiredArgsConstructor +public class SyncEventFailureService { + + private final SyncOutboxEventRepository syncOutboxEventRepository; + private final SyncEventRetrySchedule syncEventRetrySchedule; + private final Clock clock; + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public void recordFailure(UUID eventId, UUID claimToken, String errorCode, String errorMessage) { + LocalDateTime failedAt = LocalDateTime.now(clock); + SyncOutboxEvent event = syncOutboxEventRepository.findByEventIdForUpdate(eventId) + .orElseThrow(() -> new DocGridException(ErrorCode.SYNC_EVENT_NOT_FOUND)); + validateOwnership(event, claimToken, failedAt); + + // 이번 실패가 허용 횟수를 채우면 다시 Claim되지 않는 최종 상태로 종결한다. + if (event.getRetryCount() + 1 >= event.getMaxRetryCount()) { + event.markFailed(claimToken, errorCode, errorMessage, failedAt); + return; + } + event.scheduleRetry( + claimToken, + errorCode, + errorMessage, + failedAt, + syncEventRetrySchedule.nextAvailableAt(event, failedAt) + ); + } + + private void validateOwnership( + SyncOutboxEvent event, + UUID claimToken, + LocalDateTime failedAt + ) { + if (event.getStatus() != SyncEventStatus.PROCESSING + || !Objects.equals(event.getClaimToken(), claimToken) + || event.getLockExpiresAt() == null + || !event.getLockExpiresAt().isAfter(failedAt)) { + throw new DocGridException(ErrorCode.SYNC_EVENT_OWNERSHIP_INVALID); + } + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventLeaseRecoveryService.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventLeaseRecoveryService.java new file mode 100644 index 0000000..5b43841 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventLeaseRecoveryService.java @@ -0,0 +1,54 @@ +package com.opensource.docgrid.domain.sync.service.command; + +import java.time.LocalDateTime; +import java.util.Optional; +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.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncEventStatus; +import com.opensource.docgrid.domain.sync.repository.SyncOutboxEventRepository; +import com.opensource.docgrid.domain.sync.service.SyncEventRetrySchedule; + +import lombok.RequiredArgsConstructor; + +/** + * 만료 PROCESSING Event 후보를 독립 Transaction에서 다시 잠그고 Retry 또는 최종 실패로 회수한다. + */ +@Service +@RequiredArgsConstructor +public class SyncEventLeaseRecoveryService { + + private static final String LEASE_EXPIRED_CODE = "SYNC_LEASE_EXPIRED"; + private static final String LEASE_EXPIRED_MESSAGE = "Sync Dispatcher Lease가 만료되어 실행을 회수했습니다."; + + private final SyncOutboxEventRepository syncOutboxEventRepository; + private final SyncEventRetrySchedule syncEventRetrySchedule; + + @Transactional(propagation = Propagation.REQUIRES_NEW) + public RecoveryResult recover(UUID eventId, LocalDateTime recoveredAt) { + Optional candidate = syncOutboxEventRepository + .findExpiredByEventIdForUpdateSkipLocked(eventId, recoveredAt); + if (candidate.isEmpty()) { + return new RecoveryResult(eventId, false, null); + } + + SyncOutboxEvent event = candidate.get(); + event.recoverExpiredLease( + LEASE_EXPIRED_CODE, + LEASE_EXPIRED_MESSAGE, + recoveredAt, + syncEventRetrySchedule.nextAvailableAt(event, recoveredAt) + ); + return new RecoveryResult(eventId, true, event.getStatus()); + } + + /** + * 만료 후보가 실제 상태 전이를 수행했는지와 회수 후 상태를 Scheduler에 전달한다. + */ + public record RecoveryResult(UUID eventId, boolean recovered, SyncEventStatus status) { + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventLeaseService.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventLeaseService.java new file mode 100644 index 0000000..8230b61 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventLeaseService.java @@ -0,0 +1,42 @@ +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.Transactional; + +import com.opensource.docgrid.domain.sync.config.SyncDispatcherProperties; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.repository.SyncOutboxEventRepository; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +import lombok.RequiredArgsConstructor; + +/** + * 장시간 실행되는 Sync Handler가 현재 Claim Token을 증명하고 Lease만 연장할 수 있게 한다. + */ +@Service +@RequiredArgsConstructor +@Transactional +public class SyncEventLeaseService { + + private final SyncOutboxEventRepository syncOutboxEventRepository; + private final SyncDispatcherProperties syncDispatcherProperties; + private final Clock clock; + + public LocalDateTime renew(UUID eventId, UUID claimToken) { + LocalDateTime renewedAt = LocalDateTime.now(clock); + SyncOutboxEvent event = syncOutboxEventRepository.findByEventIdForUpdate(eventId) + .orElseThrow(() -> new DocGridException(ErrorCode.SYNC_EVENT_NOT_FOUND)); + LocalDateTime lockExpiresAt = renewedAt.plus(syncDispatcherProperties.getLeaseDuration()); + try { + event.renewLease(claimToken, renewedAt, lockExpiresAt); + } catch (IllegalStateException | IllegalArgumentException exception) { + throw new DocGridException(ErrorCode.SYNC_EVENT_OWNERSHIP_INVALID, exception); + } + return lockExpiresAt; + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventManualRetryService.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventManualRetryService.java new file mode 100644 index 0000000..213e648 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventManualRetryService.java @@ -0,0 +1,39 @@ +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.Transactional; + +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncEventStatus; +import com.opensource.docgrid.domain.sync.repository.SyncOutboxEventRepository; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +import lombok.RequiredArgsConstructor; + +/** + * 관리자가 최종 실패한 Sync Event에 실행 기회 한 번을 추가하는 Command Service다. + * + *

관리 API는 후속 단계에서 이 Service를 호출하며, 이 경계는 처리 중인 Event의 소유권을 덮어쓰지 않는다. + */ +@Service +@RequiredArgsConstructor +@Transactional +public class SyncEventManualRetryService { + + private final SyncOutboxEventRepository syncOutboxEventRepository; + private final Clock clock; + + public void retry(UUID eventId) { + SyncOutboxEvent event = syncOutboxEventRepository.findByEventIdForUpdate(eventId) + .orElseThrow(() -> new DocGridException(ErrorCode.SYNC_EVENT_NOT_FOUND)); + if (event.getStatus() != SyncEventStatus.FAILED) { + throw new DocGridException(ErrorCode.SYNC_EVENT_RETRY_NOT_ALLOWED); + } + event.requeueFailed(LocalDateTime.now(clock)); + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventWriter.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventWriter.java index e5c57c2..c0ee7ae 100644 --- a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventWriter.java +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventWriter.java @@ -9,9 +9,11 @@ import com.opensource.docgrid.domain.document.entity.DocumentVersion; import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel; +import com.opensource.docgrid.domain.permission.enums.AccessSourceType; import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; import com.opensource.docgrid.domain.sync.enums.SyncAggregateType; import com.opensource.docgrid.domain.sync.enums.SyncEventType; +import com.opensource.docgrid.domain.sync.enums.SyncPermissionOperation; import com.opensource.docgrid.domain.sync.repository.SyncOutboxEventRepository; import lombok.RequiredArgsConstructor; @@ -65,4 +67,90 @@ public SyncOutboxEvent recordDocumentVersionCreated( .build() ); } + + /** + * 권한 원장 변경 뒤 사용자 접근 캐시를 다시 맞출 작업을 기록한다. + */ + public SyncOutboxEvent recordPermissionCacheRefresh( + AccessSourceType sourceType, + Long sourceId, + SyncPermissionOperation operation + ) { + LocalDateTime occurredAt = LocalDateTime.now(clock); + String idempotencyKey = String.format( + "%s:%s:%d:%s:%s", + SyncAggregateType.PERMISSION, + sourceType, + sourceId, + SyncEventType.PERMISSION_CACHE_REFRESH_REQUESTED, + operation + ); + String payloadJson = String.format( + "{\"sourceType\":\"%s\",\"operation\":\"%s\"}", + sourceType, + operation + ); + return saveEvent( + idempotencyKey, + SyncAggregateType.PERMISSION, + sourceId, + null, + SyncEventType.PERMISSION_CACHE_REFRESH_REQUESTED, + payloadJson, + occurredAt + ); + } + + /** + * Reconciler나 모델 전환 흐름이 지정 Version의 재인덱싱을 멱등 요청한다. + */ + public SyncOutboxEvent recordDocumentReindexRequested( + DocumentVersion documentVersion, + EmbeddingModel embeddingModel, + String requestKey + ) { + LocalDateTime occurredAt = LocalDateTime.now(clock); + String idempotencyKey = String.format( + "%s:%d:%s:%d:%s", + SyncAggregateType.DOCUMENT_VERSION, + documentVersion.getId(), + SyncEventType.DOCUMENT_REINDEX_REQUESTED, + embeddingModel.getId(), + requestKey + ); + return saveEvent( + idempotencyKey, + SyncAggregateType.DOCUMENT_VERSION, + documentVersion.getId(), + (long) documentVersion.getVersionNo(), + SyncEventType.DOCUMENT_REINDEX_REQUESTED, + String.format("{\"embeddingModelId\":%d}", embeddingModel.getId()), + occurredAt + ); + } + + private SyncOutboxEvent saveEvent( + String idempotencyKey, + SyncAggregateType aggregateType, + Long aggregateId, + Long aggregateVersion, + SyncEventType eventType, + String payloadJson, + LocalDateTime occurredAt + ) { + return syncOutboxEventRepository.save( + SyncOutboxEvent.builder() + .eventId(UUID.randomUUID()) + .idempotencyKey(idempotencyKey) + .aggregateType(aggregateType) + .aggregateId(aggregateId) + .aggregateVersion(aggregateVersion) + .eventType(eventType) + .payloadJson(payloadJson) + .occurredAt(occurredAt) + .availableAt(occurredAt) + .maxRetryCount(DEFAULT_MAX_RETRY_COUNT) + .build() + ); + } } diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/DocumentDeletedSyncEventHandler.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/DocumentDeletedSyncEventHandler.java new file mode 100644 index 0000000..1a01326 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/DocumentDeletedSyncEventHandler.java @@ -0,0 +1,45 @@ +package com.opensource.docgrid.domain.sync.service.handler; + +import java.util.Set; + +import org.springframework.stereotype.Component; + +import com.opensource.docgrid.domain.document.entity.Document; +import com.opensource.docgrid.domain.document.enums.DocumentStatus; +import com.opensource.docgrid.domain.document.repository.DocumentRepository; +import com.opensource.docgrid.domain.embedding.repository.EmbeddingRepository; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncEventType; +import com.opensource.docgrid.domain.sync.service.SyncEventHandler; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +import lombok.RequiredArgsConstructor; + +/** + * Soft-delete된 문서의 ACTIVE Vector를 STALE로 전환해 검색에서 즉시 제외한다. + * + *

문서 원장 상태가 DELETED인지 다시 확인하며 Chunk와 Vector의 물리 삭제는 수행하지 않는다. + */ +@Component +@RequiredArgsConstructor +public class DocumentDeletedSyncEventHandler implements SyncEventHandler { + + private final DocumentRepository documentRepository; + private final EmbeddingRepository embeddingRepository; + + @Override + public Set supportedTypes() { + return Set.of(SyncEventType.DOCUMENT_DELETED); + } + + @Override + public void handle(SyncOutboxEvent event) { + Document document = documentRepository.findById(event.getAggregateId()) + .orElseThrow(() -> new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT)); + if (document.getStatus() != DocumentStatus.DELETED || document.getDeletedAt() == null) { + throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT); + } + embeddingRepository.markActiveAsStaleByDocumentId(document.getId()); + } +} 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 new file mode 100644 index 0000000..00fde35 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/DocumentVersionSyncEventHandler.java @@ -0,0 +1,114 @@ +package com.opensource.docgrid.domain.sync.service.handler; + +import java.util.EnumSet; +import java.util.Objects; +import java.util.Optional; +import java.util.Set; + +import org.springframework.stereotype.Component; + +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.EmbeddingJob; +import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel; +import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus; +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.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncEventType; +import com.opensource.docgrid.domain.sync.service.SyncEventHandler; +import com.opensource.docgrid.domain.sync.service.SyncEventPayloadReader; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +import lombok.RequiredArgsConstructor; + +/** + * 문서 버전 생성과 명시적 재인덱싱 Event를 기존 Embedding Job Queue에 멱등 반영한다. + * + *

최초 업로드 Transaction에서 이미 같은 sourceEventId Job이 생성됐다면 검증만 하고, 과거 장애로 + * Job이 누락된 경우에만 새 Job을 만든다. 실패 Job 재인덱싱은 기존 수동 재처리 불변식을 재사용한다. + */ +@Component +@RequiredArgsConstructor +public class DocumentVersionSyncEventHandler implements SyncEventHandler { + + private static final int DEFAULT_JOB_PRIORITY = 0; + private static final int MAX_RETRY_COUNT = 3; + private static final Set SUPPORTED_TYPES = EnumSet.of( + SyncEventType.DOCUMENT_VERSION_CREATED, + SyncEventType.DOCUMENT_REINDEX_REQUESTED + ); + + private final DocumentVersionRepository documentVersionRepository; + private final EmbeddingModelRepository embeddingModelRepository; + private final EmbeddingJobRepository embeddingJobRepository; + private final EmbeddingJobManualRetryService embeddingJobManualRetryService; + private final SyncEventPayloadReader payloadReader; + + @Override + public Set supportedTypes() { + return SUPPORTED_TYPES; + } + + @Override + public void handle(SyncOutboxEvent event) { + DocumentVersion version = documentVersionRepository.findById(event.getAggregateId()) + .orElseThrow(() -> new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT)); + EmbeddingModel model = embeddingModelRepository.findById( + payloadReader.requiredLong(event, "embeddingModelId") + ) + .orElseThrow(() -> new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT)); + + // 같은 sourceEventId Job은 최초 실행이나 재전달 모두 검증 후 그대로 재사용한다. + Optional sourceJob = embeddingJobRepository.findBySourceEventId(event.getEventId()); + if (sourceJob.isPresent()) { + validateJob(sourceJob.get(), version, model); + return; + } + + Optional latestJob = embeddingJobRepository + .findTopByDocumentVersionIdAndEmbeddingModelIdOrderByIdDesc(version.getId(), model.getId()); + if (event.getEventType() == SyncEventType.DOCUMENT_REINDEX_REQUESTED && latestJob.isPresent()) { + retryOrReuse(latestJob.get()); + return; + } + if (latestJob.isPresent() || version.getStatus() != DocumentVersionStatus.UPLOADED) { + throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT); + } + + embeddingJobRepository.save( + EmbeddingJob.builder() + .documentVersion(version) + .embeddingModel(model) + .sourceEventId(event.getEventId()) + .status(EmbeddingJobStatus.PENDING) + .priority(DEFAULT_JOB_PRIORITY) + .maxRetryCount(MAX_RETRY_COUNT) + .build() + ); + } + + private void retryOrReuse(EmbeddingJob job) { + if (job.getStatus() == EmbeddingJobStatus.FAILED) { + embeddingJobManualRetryService.retry(job.getId()); + return; + } + if (job.getStatus() == EmbeddingJobStatus.PENDING + || job.getStatus() == EmbeddingJobStatus.PROCESSING) { + return; + } + throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT); + } + + private void validateJob(EmbeddingJob job, DocumentVersion version, EmbeddingModel model) { + if (job.getDocumentVersion() == null + || job.getEmbeddingModel() == null + || !Objects.equals(job.getDocumentVersion().getId(), version.getId()) + || !Objects.equals(job.getEmbeddingModel().getId(), model.getId())) { + throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT); + } + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/EmbeddingModelActivatedSyncEventHandler.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/EmbeddingModelActivatedSyncEventHandler.java new file mode 100644 index 0000000..570c480 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/EmbeddingModelActivatedSyncEventHandler.java @@ -0,0 +1,42 @@ +package com.opensource.docgrid.domain.sync.service.handler; + +import java.util.Set; + +import org.springframework.stereotype.Component; + +import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel; +import com.opensource.docgrid.domain.embedding.repository.EmbeddingModelRepository; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncEventType; +import com.opensource.docgrid.domain.sync.service.SyncEventHandler; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +import lombok.RequiredArgsConstructor; + +/** + * 모델 활성화 Event가 현재 검색 모델 원장과 일치하는지 검증한다. + * + *

대량 문서 재인덱싱 대상 산출과 안전한 Repair Event 생성은 Reconciler 단계가 담당한다. 이 Handler는 + * 비활성 모델 Event가 성공 처리되어 잘못된 기준점이 남는 것을 차단한다. + */ +@Component +@RequiredArgsConstructor +public class EmbeddingModelActivatedSyncEventHandler implements SyncEventHandler { + + private final EmbeddingModelRepository embeddingModelRepository; + + @Override + public Set supportedTypes() { + return Set.of(SyncEventType.EMBEDDING_MODEL_ACTIVATED); + } + + @Override + public void handle(SyncOutboxEvent event) { + EmbeddingModel model = embeddingModelRepository.findById(event.getAggregateId()) + .orElseThrow(() -> new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT)); + if (!model.isActive() || !model.isSearchable()) { + throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT); + } + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/PermissionCacheRefreshSyncEventHandler.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/PermissionCacheRefreshSyncEventHandler.java new file mode 100644 index 0000000..efaf676 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/handler/PermissionCacheRefreshSyncEventHandler.java @@ -0,0 +1,127 @@ +package com.opensource.docgrid.domain.sync.service.handler; + +import java.util.List; +import java.util.Optional; +import java.util.Set; + +import org.springframework.stereotype.Component; + +import com.opensource.docgrid.domain.collection.entity.CollectionDocument; +import com.opensource.docgrid.domain.collection.repository.CollectionDocumentRepository; +import com.opensource.docgrid.domain.document.entity.Document; +import com.opensource.docgrid.domain.permission.entity.CollectionPermission; +import com.opensource.docgrid.domain.permission.entity.DocumentPermission; +import com.opensource.docgrid.domain.permission.enums.AccessSourceType; +import com.opensource.docgrid.domain.permission.enums.PermissionTargetType; +import com.opensource.docgrid.domain.permission.repository.CollectionPermissionRepository; +import com.opensource.docgrid.domain.permission.repository.DocumentPermissionRepository; +import com.opensource.docgrid.domain.permission.service.command.UserDocumentAccessCacheService; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncEventType; +import com.opensource.docgrid.domain.sync.enums.SyncPermissionOperation; +import com.opensource.docgrid.domain.sync.service.SyncEventHandler; +import com.opensource.docgrid.domain.sync.service.SyncEventPayloadReader; +import com.opensource.docgrid.global.exception.DocGridException; +import com.opensource.docgrid.global.exception.ErrorCode; + +import lombok.RequiredArgsConstructor; + +/** + * 직접 USER 문서·컬렉션 권한 원장을 접근 캐시에 다시 투영한다. + * + *

권한이 이미 회수됐거나 ROLE·DEPARTMENT 대상이면 해당 출처 캐시를 무효화한다. 캐시는 최종 권한 + * 원장이 아니므로 이 Handler는 Vector를 변경하지 않는다. + */ +@Component +@RequiredArgsConstructor +public class PermissionCacheRefreshSyncEventHandler implements SyncEventHandler { + + private final DocumentPermissionRepository documentPermissionRepository; + private final CollectionPermissionRepository collectionPermissionRepository; + private final CollectionDocumentRepository collectionDocumentRepository; + private final UserDocumentAccessCacheService cacheService; + private final SyncEventPayloadReader payloadReader; + + @Override + public Set supportedTypes() { + return Set.of(SyncEventType.PERMISSION_CACHE_REFRESH_REQUESTED); + } + + @Override + public void handle(SyncOutboxEvent event) { + AccessSourceType sourceType = parseSourceType(event); + SyncPermissionOperation operation = parseOperation(event); + if (event.getAggregateId() == null || sourceType == AccessSourceType.OWNER) { + throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT); + } + if (operation == SyncPermissionOperation.REVOKED) { + cacheService.bulkRevokeBySource(sourceType, event.getAggregateId()); + return; + } + + switch (sourceType) { + case DIRECT_DOCUMENT_PERMISSION -> refreshDocumentPermission(event.getAggregateId()); + case DIRECT_COLLECTION_PERMISSION -> refreshCollectionPermission(event.getAggregateId()); + case OWNER -> throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT); + } + } + + private void refreshDocumentPermission(Long permissionId) { + Optional permission = documentPermissionRepository.findById(permissionId); + if (permission.isEmpty() || permission.get().getTargetType() != PermissionTargetType.USER) { + cacheService.bulkRevokeBySource(AccessSourceType.DIRECT_DOCUMENT_PERMISSION, permissionId); + return; + } + DocumentPermission current = permission.get(); + cacheService.grantUserPermission( + current.getUser(), + current.getDocument(), + current.isCanRead(), + current.isCanWrite(), + current.isCanAdmin(), + AccessSourceType.DIRECT_DOCUMENT_PERMISSION, + permissionId, + current.getExpiresAt() + ); + } + + private void refreshCollectionPermission(Long permissionId) { + Optional permission = collectionPermissionRepository.findById(permissionId); + if (permission.isEmpty() || permission.get().getTargetType() != PermissionTargetType.USER) { + cacheService.bulkRevokeBySource(AccessSourceType.DIRECT_COLLECTION_PERMISSION, permissionId); + return; + } + CollectionPermission current = permission.get(); + List documents = collectionDocumentRepository + .findAllByCollectionId(current.getCollection().getId()) + .stream() + .map(CollectionDocument::getDocument) + .toList(); + cacheService.bulkGrantUserPermission( + current.getUser(), + documents, + current.isCanRead(), + current.isCanWrite(), + current.isCanAdmin(), + AccessSourceType.DIRECT_COLLECTION_PERMISSION, + permissionId, + current.getExpiresAt() + ); + } + + private AccessSourceType parseSourceType(SyncOutboxEvent event) { + try { + return AccessSourceType.valueOf(payloadReader.requiredText(event, "sourceType")); + } catch (IllegalArgumentException exception) { + throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT, exception); + } + } + + private SyncPermissionOperation parseOperation(SyncOutboxEvent event) { + try { + return SyncPermissionOperation.valueOf(payloadReader.requiredText(event, "operation")); + } catch (IllegalArgumentException exception) { + throw new DocGridException(ErrorCode.SYNC_EVENT_INCONSISTENT, exception); + } + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/global/exception/ErrorCode.java b/backend/src/main/java/com/opensource/docgrid/global/exception/ErrorCode.java index 74c9732..499166e 100644 --- a/backend/src/main/java/com/opensource/docgrid/global/exception/ErrorCode.java +++ b/backend/src/main/java/com/opensource/docgrid/global/exception/ErrorCode.java @@ -240,6 +240,28 @@ public enum ErrorCode { "사용 가능한 임베딩 모델이 여러 개 설정되어 있습니다." ), + // SYNC + SYNC_EVENT_NOT_FOUND( + HttpStatus.NOT_FOUND, + "SYNC-001", + "동기화 Event를 찾을 수 없습니다." + ), + SYNC_EVENT_OWNERSHIP_INVALID( + HttpStatus.CONFLICT, + "SYNC-002", + "현재 동기화 Event 소유권과 요청이 일치하지 않습니다." + ), + SYNC_EVENT_INCONSISTENT( + HttpStatus.INTERNAL_SERVER_ERROR, + "SYNC-003", + "동기화 Event와 도메인 상태가 일치하지 않습니다." + ), + SYNC_EVENT_RETRY_NOT_ALLOWED( + HttpStatus.CONFLICT, + "SYNC-004", + "최종 실패한 동기화 Event만 재처리할 수 있습니다." + ), + // SEARCH EMBEDDING_SERVER_UNAVAILABLE( HttpStatus.SERVICE_UNAVAILABLE, diff --git a/backend/src/main/resources/application.yml b/backend/src/main/resources/application.yml index 48063fb..1118713 100644 --- a/backend/src/main/resources/application.yml +++ b/backend/src/main/resources/application.yml @@ -57,6 +57,17 @@ indexing: retry-max-delay: ${INDEXING_WORKER_RETRY_MAX_DELAY:5m} shutdown-grace-period: ${INDEXING_WORKER_SHUTDOWN_GRACE_PERIOD:30s} +sync: + dispatcher: + enabled: ${SYNC_DISPATCHER_ENABLED:false} + name: ${SYNC_DISPATCHER_NAME:sync-dispatcher} + polling-interval: ${SYNC_DISPATCHER_POLLING_INTERVAL:1s} + lease-duration: ${SYNC_DISPATCHER_LEASE_DURATION:30s} + lease-recovery-interval: ${SYNC_DISPATCHER_LEASE_RECOVERY_INTERVAL:10s} + 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} + server: port: 8080 diff --git a/backend/src/test/java/com/opensource/docgrid/domain/permission/service/command/CollectionPermissionCommandServiceTest.java b/backend/src/test/java/com/opensource/docgrid/domain/permission/service/command/CollectionPermissionCommandServiceTest.java index 188a09c..8ad004c 100644 --- a/backend/src/test/java/com/opensource/docgrid/domain/permission/service/command/CollectionPermissionCommandServiceTest.java +++ b/backend/src/test/java/com/opensource/docgrid/domain/permission/service/command/CollectionPermissionCommandServiceTest.java @@ -37,6 +37,7 @@ import com.opensource.docgrid.domain.user.repository.DepartmentRepository; import com.opensource.docgrid.domain.user.repository.RoleRepository; import com.opensource.docgrid.domain.user.repository.UserRepository; +import com.opensource.docgrid.domain.sync.service.command.SyncEventWriter; import com.opensource.docgrid.global.exception.DocGridException; import com.opensource.docgrid.global.exception.ErrorCode; @@ -56,6 +57,7 @@ class CollectionPermissionCommandServiceTest { @Mock private DepartmentRepository departmentRepository; @Mock private PermissionConverter permissionConverter; @Mock private PermissionQueryService permissionQueryService; + @Mock private SyncEventWriter syncEventWriter; // ==================== grantPermission ==================== diff --git a/backend/src/test/java/com/opensource/docgrid/domain/permission/service/command/DocumentPermissionCommandServiceTest.java b/backend/src/test/java/com/opensource/docgrid/domain/permission/service/command/DocumentPermissionCommandServiceTest.java index 3dd10de..6695a27 100644 --- a/backend/src/test/java/com/opensource/docgrid/domain/permission/service/command/DocumentPermissionCommandServiceTest.java +++ b/backend/src/test/java/com/opensource/docgrid/domain/permission/service/command/DocumentPermissionCommandServiceTest.java @@ -30,6 +30,7 @@ import com.opensource.docgrid.domain.user.repository.DepartmentRepository; import com.opensource.docgrid.domain.user.repository.RoleRepository; import com.opensource.docgrid.domain.user.repository.UserRepository; +import com.opensource.docgrid.domain.sync.service.command.SyncEventWriter; import com.opensource.docgrid.global.exception.DocGridException; import com.opensource.docgrid.global.exception.ErrorCode; @@ -48,6 +49,7 @@ class DocumentPermissionCommandServiceTest { @Mock private DepartmentRepository departmentRepository; @Mock private PermissionConverter permissionConverter; @Mock private PermissionQueryService permissionQueryService; + @Mock private SyncEventWriter syncEventWriter; @Test @DisplayName("USER 대상 문서 권한을 부여하면 캐시도 함께 갱신된다") diff --git a/backend/src/test/java/com/opensource/docgrid/domain/sync/entity/SyncOutboxEventTest.java b/backend/src/test/java/com/opensource/docgrid/domain/sync/entity/SyncOutboxEventTest.java new file mode 100644 index 0000000..47b6501 --- /dev/null +++ b/backend/src/test/java/com/opensource/docgrid/domain/sync/entity/SyncOutboxEventTest.java @@ -0,0 +1,83 @@ +package com.opensource.docgrid.domain.sync.entity; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.time.LocalDateTime; +import java.util.UUID; + +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Test; + +import com.opensource.docgrid.domain.sync.enums.SyncAggregateType; +import com.opensource.docgrid.domain.sync.enums.SyncEventStatus; +import com.opensource.docgrid.domain.sync.enums.SyncEventType; + +/** + * Sync Outbox Event의 Claim, 완료, Retry와 만료 Lease 복구 상태 불변식을 검증한다. + */ +@DisplayName("SyncOutboxEvent 상태 전이 테스트") +class SyncOutboxEventTest { + + private static final LocalDateTime NOW = LocalDateTime.of(2026, 8, 13, 18, 30); + + @Test + @DisplayName("PENDING Event를 Claim하고 같은 Token으로 완료한다") + void claimAndComplete_transitionsToProcessed() { + SyncOutboxEvent event = event(3); + UUID claimToken = UUID.randomUUID(); + + event.claim("dispatcher-1", claimToken, NOW, NOW.plusSeconds(30)); + event.complete(claimToken, NOW.plusSeconds(1)); + + assertThat(event.getStatus()).isEqualTo(SyncEventStatus.PROCESSED); + assertThat(event.getProcessedAt()).isEqualTo(NOW.plusSeconds(1)); + assertThat(event.getClaimToken()).isNull(); + assertThat(event.getLockedBy()).isNull(); + assertThat(event.getLockExpiresAt()).isNull(); + } + + @Test + @DisplayName("다른 Claim Token으로 완료할 수 없다") + void completeRejectsStaleClaimToken() { + SyncOutboxEvent event = event(3); + event.claim("dispatcher-1", UUID.randomUUID(), NOW, NOW.plusSeconds(30)); + + assertThatThrownBy(() -> event.complete(UUID.randomUUID(), NOW.plusSeconds(1))) + .isInstanceOf(IllegalStateException.class); + assertThat(event.getStatus()).isEqualTo(SyncEventStatus.PROCESSING); + } + + @Test + @DisplayName("만료 Lease는 Retry Queue로 회수하고 횟수를 소진하면 FAILED로 종결한다") + void recoverExpiredLease_retriesThenFails() { + SyncOutboxEvent event = event(2); + event.claim("dispatcher-1", UUID.randomUUID(), NOW, NOW.plusSeconds(10)); + + event.recoverExpiredLease("EXPIRED", "expired", NOW.plusSeconds(10), NOW.plusSeconds(15)); + assertThat(event.getStatus()).isEqualTo(SyncEventStatus.PENDING); + assertThat(event.getRetryCount()).isOne(); + + event.claim("dispatcher-2", UUID.randomUUID(), NOW.plusSeconds(15), NOW.plusSeconds(25)); + event.recoverExpiredLease("EXPIRED", "expired", NOW.plusSeconds(25), NOW.plusSeconds(30)); + + assertThat(event.getStatus()).isEqualTo(SyncEventStatus.FAILED); + assertThat(event.getRetryCount()).isEqualTo(2); + assertThat(event.getClaimToken()).isNull(); + } + + private SyncOutboxEvent event(int maxRetryCount) { + return SyncOutboxEvent.builder() + .eventId(UUID.randomUUID()) + .idempotencyKey("test:" + UUID.randomUUID()) + .aggregateType(SyncAggregateType.DOCUMENT_VERSION) + .aggregateId(1L) + .aggregateVersion(1L) + .eventType(SyncEventType.DOCUMENT_VERSION_CREATED) + .payloadJson("{\"embeddingModelId\":1}") + .availableAt(NOW) + .occurredAt(NOW) + .maxRetryCount(maxRetryCount) + .build(); + } +} diff --git a/backend/src/test/java/com/opensource/docgrid/domain/sync/integration/SyncEventClaimIntegrationTest.java b/backend/src/test/java/com/opensource/docgrid/domain/sync/integration/SyncEventClaimIntegrationTest.java new file mode 100644 index 0000000..6743874 --- /dev/null +++ b/backend/src/test/java/com/opensource/docgrid/domain/sync/integration/SyncEventClaimIntegrationTest.java @@ -0,0 +1,114 @@ +package com.opensource.docgrid.domain.sync.integration; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.time.LocalDateTime; +import java.util.ArrayList; +import java.util.List; +import java.util.Optional; +import java.util.UUID; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.DisplayName; +import org.junit.jupiter.api.Tag; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.test.context.ActiveProfiles; + +import com.opensource.docgrid.domain.sync.dto.ClaimedSyncEvent; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncAggregateType; +import com.opensource.docgrid.domain.sync.enums.SyncEventStatus; +import com.opensource.docgrid.domain.sync.enums.SyncEventType; +import com.opensource.docgrid.domain.sync.repository.SyncOutboxEventRepository; +import com.opensource.docgrid.domain.sync.service.command.SyncEventClaimService; + +/** + * 실제 PostgreSQL SKIP LOCKED Queue에서 다중 Dispatcher가 한 Event를 중복 Claim하지 않는지 검증한다. + */ +@Tag("integration") +@SpringBootTest +@ActiveProfiles("test") +@DisplayName("Sync Event Claim 통합 테스트") +class SyncEventClaimIntegrationTest { + + private static final int COMPETITOR_COUNT = 20; + + @Autowired private SyncOutboxEventRepository syncOutboxEventRepository; + @Autowired private SyncEventClaimService syncEventClaimService; + @Autowired private JdbcTemplate jdbcTemplate; + + private UUID eventId; + + @BeforeEach + void setUp() { + // 이 Test Schema에 남은 Event가 경쟁 대상에 섞이지 않도록 기존 Queue를 종결 상태로 격리한다. + jdbcTemplate.update(""" + UPDATE sync_outbox_events + SET status = 'PROCESSED', processed_at = CURRENT_TIMESTAMP, + claim_token = NULL, locked_by = NULL, lock_expires_at = NULL + WHERE status IN ('PENDING', 'PROCESSING') + """); + LocalDateTime now = LocalDateTime.now().minusSeconds(1); + SyncOutboxEvent event = syncOutboxEventRepository.saveAndFlush( + SyncOutboxEvent.builder() + .eventId(UUID.randomUUID()) + .idempotencyKey("claim-integration:" + UUID.randomUUID()) + .aggregateType(SyncAggregateType.DOCUMENT_VERSION) + .aggregateId(9_999_999L) + .aggregateVersion(1L) + .eventType(SyncEventType.DOCUMENT_VERSION_CREATED) + .payloadJson("{\"embeddingModelId\":1}") + .availableAt(now) + .occurredAt(now) + .maxRetryCount(5) + .build() + ); + eventId = event.getEventId(); + } + + @AfterEach + void tearDown() { + syncOutboxEventRepository.findByEventId(eventId) + .ifPresent(syncOutboxEventRepository::delete); + } + + @Test + @DisplayName("20개 Dispatcher가 동시에 경쟁해도 같은 Event Claim은 한 건이다") + void claim_assignsSingleOwner_whenTwentyDispatchersCompete() throws Exception { + CyclicBarrier barrier = new CyclicBarrier(COMPETITOR_COUNT); + ExecutorService executor = Executors.newFixedThreadPool(COMPETITOR_COUNT); + List>> futures = new ArrayList<>(); + try { + for (int index = 0; index < COMPETITOR_COUNT; index++) { + futures.add(executor.submit(() -> { + barrier.await(10, TimeUnit.SECONDS); + return syncEventClaimService.claim(); + })); + } + + List claims = new ArrayList<>(); + for (Future> future : futures) { + future.get(15, TimeUnit.SECONDS).ifPresent(claims::add); + } + + assertThat(claims).hasSize(1); + assertThat(claims.get(0).eventId()).isEqualTo(eventId); + assertThat(claims.get(0).claimToken()).isNotNull(); + SyncOutboxEvent persisted = syncOutboxEventRepository.findByEventId(eventId).orElseThrow(); + assertThat(persisted.getStatus()).isEqualTo(SyncEventStatus.PROCESSING); + assertThat(persisted.getClaimToken()).isEqualTo(claims.get(0).claimToken()); + } finally { + executor.shutdownNow(); + assertThat(executor.awaitTermination(10, TimeUnit.SECONDS)).isTrue(); + } + } +} diff --git a/backend/src/test/java/com/opensource/docgrid/domain/sync/service/command/SyncEventClaimServiceTest.java b/backend/src/test/java/com/opensource/docgrid/domain/sync/service/command/SyncEventClaimServiceTest.java new file mode 100644 index 0000000..cd68d4e --- /dev/null +++ b/backend/src/test/java/com/opensource/docgrid/domain/sync/service/command/SyncEventClaimServiceTest.java @@ -0,0 +1,87 @@ +package com.opensource.docgrid.domain.sync.service.command; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.BDDMockito.given; + +import java.time.Clock; +import java.time.Duration; +import java.time.Instant; +import java.time.LocalDateTime; +import java.time.ZoneId; +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.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import com.opensource.docgrid.domain.sync.config.SyncDispatcherProperties; +import com.opensource.docgrid.domain.sync.dto.ClaimedSyncEvent; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncAggregateType; +import com.opensource.docgrid.domain.sync.enums.SyncEventStatus; +import com.opensource.docgrid.domain.sync.enums.SyncEventType; +import com.opensource.docgrid.domain.sync.repository.SyncOutboxEventRepository; + +/** + * SyncEventClaimService가 Queue 조회와 Lease 발급을 하나의 상태 전이로 수행하는지 검증한다. + */ +@ExtendWith(MockitoExtension.class) +@DisplayName("SyncEventClaimService 단위 테스트") +class SyncEventClaimServiceTest { + + private static final Instant NOW = Instant.parse("2026-08-13T09:30:00Z"); + private static final ZoneId ZONE_ID = ZoneId.of("Asia/Seoul"); + + @Mock private SyncOutboxEventRepository syncOutboxEventRepository; + + private SyncEventClaimService service; + private SyncDispatcherProperties properties; + + @BeforeEach + void setUp() { + properties = new SyncDispatcherProperties(); + properties.setName("dispatcher-test"); + properties.setLeaseDuration(Duration.ofSeconds(30)); + service = new SyncEventClaimService( + syncOutboxEventRepository, + properties, + Clock.fixed(NOW, ZONE_ID) + ); + } + + @Test + @DisplayName("가장 오래된 PENDING Event에 고유 Claim Token과 Lease를 부여한다") + void claim_assignsOwnershipAndLease() { + LocalDateTime claimedAt = LocalDateTime.ofInstant(NOW, ZONE_ID); + SyncOutboxEvent event = event(claimedAt); + given(syncOutboxEventRepository.findNextPendingForUpdate(claimedAt)).willReturn(Optional.of(event)); + + Optional result = service.claim(); + + assertThat(result).isPresent(); + assertThat(result.orElseThrow().eventId()).isEqualTo(event.getEventId()); + assertThat(result.orElseThrow().claimToken()).isNotNull(); + assertThat(event.getStatus()).isEqualTo(SyncEventStatus.PROCESSING); + assertThat(event.getLockedBy()).isEqualTo("dispatcher-test"); + assertThat(event.getLockExpiresAt()).isEqualTo(claimedAt.plusSeconds(30)); + } + + private SyncOutboxEvent event(LocalDateTime availableAt) { + return SyncOutboxEvent.builder() + .eventId(UUID.randomUUID()) + .idempotencyKey("claim:" + UUID.randomUUID()) + .aggregateType(SyncAggregateType.DOCUMENT_VERSION) + .aggregateId(1L) + .aggregateVersion(1L) + .eventType(SyncEventType.DOCUMENT_VERSION_CREATED) + .payloadJson("{\"embeddingModelId\":1}") + .availableAt(availableAt) + .occurredAt(availableAt) + .maxRetryCount(3) + .build(); + } +} diff --git a/backend/src/test/java/com/opensource/docgrid/domain/sync/service/command/SyncEventDispatchServiceTest.java b/backend/src/test/java/com/opensource/docgrid/domain/sync/service/command/SyncEventDispatchServiceTest.java new file mode 100644 index 0000000..b8839f7 --- /dev/null +++ b/backend/src/test/java/com/opensource/docgrid/domain/sync/service/command/SyncEventDispatchServiceTest.java @@ -0,0 +1,97 @@ +package com.opensource.docgrid.domain.sync.service.command; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.then; +import static org.mockito.Mockito.doThrow; + +import java.time.Clock; +import java.time.Instant; +import java.time.LocalDateTime; +import java.time.ZoneId; +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.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import com.opensource.docgrid.domain.sync.dto.ClaimedSyncEvent; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncAggregateType; +import com.opensource.docgrid.domain.sync.enums.SyncEventStatus; +import com.opensource.docgrid.domain.sync.enums.SyncEventType; +import com.opensource.docgrid.domain.sync.repository.SyncOutboxEventRepository; +import com.opensource.docgrid.domain.sync.service.SyncEventHandlerRegistry; + +/** + * Handler 실행과 Outbox Event 완료 전이가 같은 Dispatch 경계에서 수행되는지 검증한다. + */ +@ExtendWith(MockitoExtension.class) +@DisplayName("SyncEventDispatchService 단위 테스트") +class SyncEventDispatchServiceTest { + + private static final Instant NOW = Instant.parse("2026-08-13T10:00:00Z"); + private static final ZoneId ZONE_ID = ZoneId.of("Asia/Seoul"); + + @Mock private SyncOutboxEventRepository syncOutboxEventRepository; + @Mock private SyncEventHandlerRegistry syncEventHandlerRegistry; + + private SyncEventDispatchService service; + private SyncOutboxEvent event; + private ClaimedSyncEvent claim; + + @BeforeEach + void setUp() { + Clock clock = Clock.fixed(NOW, ZONE_ID); + service = new SyncEventDispatchService(syncOutboxEventRepository, syncEventHandlerRegistry, clock); + LocalDateTime now = LocalDateTime.ofInstant(NOW, ZONE_ID); + event = event(now); + UUID claimToken = UUID.randomUUID(); + event.claim("dispatcher", claimToken, now, now.plusSeconds(30)); + claim = new ClaimedSyncEvent(event.getEventId(), claimToken); + given(syncOutboxEventRepository.findByEventIdForUpdate(event.getEventId())) + .willReturn(Optional.of(event)); + } + + @Test + @DisplayName("Handler 성공 후 Event를 PROCESSED로 완료한다") + void dispatch_completesEvent_afterHandlerSucceeds() { + service.dispatch(claim); + + then(syncEventHandlerRegistry).should().handle(event); + assertThat(event.getStatus()).isEqualTo(SyncEventStatus.PROCESSED); + assertThat(event.getProcessedAt()).isNotNull(); + } + + @Test + @DisplayName("Handler가 실패하면 완료 전이를 실행하지 않는다") + void dispatch_doesNotComplete_whenHandlerFails() { + doThrow(new IllegalStateException("handler failure")) + .when(syncEventHandlerRegistry).handle(event); + + assertThatThrownBy(() -> service.dispatch(claim)) + .isInstanceOf(IllegalStateException.class); + assertThat(event.getStatus()).isEqualTo(SyncEventStatus.PROCESSING); + assertThat(event.getProcessedAt()).isNull(); + } + + private SyncOutboxEvent event(LocalDateTime now) { + return SyncOutboxEvent.builder() + .eventId(UUID.randomUUID()) + .idempotencyKey("dispatch:" + UUID.randomUUID()) + .aggregateType(SyncAggregateType.DOCUMENT_VERSION) + .aggregateId(1L) + .aggregateVersion(1L) + .eventType(SyncEventType.DOCUMENT_VERSION_CREATED) + .payloadJson("{\"embeddingModelId\":1}") + .availableAt(now) + .occurredAt(now) + .maxRetryCount(3) + .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 new file mode 100644 index 0000000..f6f36ce --- /dev/null +++ b/backend/src/test/java/com/opensource/docgrid/domain/sync/service/handler/DocumentVersionSyncEventHandlerTest.java @@ -0,0 +1,124 @@ +package com.opensource.docgrid.domain.sync.service.handler; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.BDDMockito.given; +import static org.mockito.BDDMockito.then; +import static org.mockito.Mockito.never; + +import java.time.LocalDateTime; +import 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.fasterxml.jackson.databind.ObjectMapper; +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.EmbeddingJob; +import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel; +import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus; +import com.opensource.docgrid.domain.embedding.fixture.EmbeddingModelFixture; +import com.opensource.docgrid.domain.embedding.repository.EmbeddingJobRepository; +import com.opensource.docgrid.domain.embedding.repository.EmbeddingModelRepository; +import com.opensource.docgrid.domain.embedding.service.command.EmbeddingJobManualRetryService; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.enums.SyncAggregateType; +import com.opensource.docgrid.domain.sync.enums.SyncEventType; +import com.opensource.docgrid.domain.sync.service.SyncEventPayloadReader; + +/** + * 문서 버전 Event 재전달이 기존 Job을 재사용하고 누락된 경우에만 한 Job을 만드는지 검증한다. + */ +@ExtendWith(MockitoExtension.class) +@DisplayName("DocumentVersionSyncEventHandler 단위 테스트") +class DocumentVersionSyncEventHandlerTest { + + @Mock private DocumentVersionRepository documentVersionRepository; + @Mock private EmbeddingModelRepository embeddingModelRepository; + @Mock private EmbeddingJobRepository embeddingJobRepository; + @Mock private EmbeddingJobManualRetryService embeddingJobManualRetryService; + + private DocumentVersionSyncEventHandler handler; + private DocumentVersion version; + private EmbeddingModel model; + private SyncOutboxEvent event; + + @BeforeEach + void setUp() { + handler = new DocumentVersionSyncEventHandler( + documentVersionRepository, + embeddingModelRepository, + embeddingJobRepository, + embeddingJobManualRetryService, + new SyncEventPayloadReader(new ObjectMapper()) + ); + version = DocumentVersion.builder().versionNo(1).status(DocumentVersionStatus.UPLOADED).build(); + ReflectionTestUtils.setField(version, "id", 11L); + model = EmbeddingModelFixture.createDefaultModel(); + ReflectionTestUtils.setField(model, "id", 7L); + event = event(); + given(documentVersionRepository.findById(11L)).willReturn(Optional.of(version)); + given(embeddingModelRepository.findById(7L)).willReturn(Optional.of(model)); + } + + @Test + @DisplayName("같은 source Event Job이 있으면 새 Job을 만들지 않는다") + void handle_reusesExistingSourceJob() { + EmbeddingJob existing = EmbeddingJob.builder() + .documentVersion(version) + .embeddingModel(model) + .sourceEventId(event.getEventId()) + .status(EmbeddingJobStatus.PENDING) + .priority(0) + .maxRetryCount(3) + .build(); + given(embeddingJobRepository.findBySourceEventId(event.getEventId())) + .willReturn(Optional.of(existing)); + + handler.handle(event); + + then(embeddingJobRepository).should(never()).save(any()); + } + + @Test + @DisplayName("Job이 누락됐으면 source Event와 연결된 PENDING Job 하나를 생성한다") + void handle_createsJob_whenSourceJobIsMissing() { + given(embeddingJobRepository.findBySourceEventId(event.getEventId())).willReturn(Optional.empty()); + given(embeddingJobRepository.findTopByDocumentVersionIdAndEmbeddingModelIdOrderByIdDesc(11L, 7L)) + .willReturn(Optional.empty()); + + handler.handle(event); + + ArgumentCaptor jobCaptor = ArgumentCaptor.forClass(EmbeddingJob.class); + then(embeddingJobRepository).should().save(jobCaptor.capture()); + assertThat(jobCaptor.getValue().getSourceEventId()).isEqualTo(event.getEventId()); + assertThat(jobCaptor.getValue().getStatus()).isEqualTo(EmbeddingJobStatus.PENDING); + assertThat(jobCaptor.getValue().getDocumentVersion()).isSameAs(version); + assertThat(jobCaptor.getValue().getEmbeddingModel()).isSameAs(model); + } + + private SyncOutboxEvent event() { + LocalDateTime now = LocalDateTime.of(2026, 8, 13, 19, 0); + return SyncOutboxEvent.builder() + .eventId(UUID.randomUUID()) + .idempotencyKey("version-created:11") + .aggregateType(SyncAggregateType.DOCUMENT_VERSION) + .aggregateId(11L) + .aggregateVersion(1L) + .eventType(SyncEventType.DOCUMENT_VERSION_CREATED) + .payloadJson("{\"embeddingModelId\":7}") + .availableAt(now) + .occurredAt(now) + .maxRetryCount(5) + .build(); + } +}