diff --git a/backend/src/main/java/com/opensource/docgrid/domain/document/service/command/DocumentUploadService.java b/backend/src/main/java/com/opensource/docgrid/domain/document/service/command/DocumentUploadService.java index 42d1da0..67c705d 100644 --- a/backend/src/main/java/com/opensource/docgrid/domain/document/service/command/DocumentUploadService.java +++ b/backend/src/main/java/com/opensource/docgrid/domain/document/service/command/DocumentUploadService.java @@ -18,6 +18,8 @@ import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus; import com.opensource.docgrid.domain.embedding.repository.EmbeddingJobRepository; import com.opensource.docgrid.domain.embedding.service.query.EmbeddingModelQueryService; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.service.command.SyncEventWriter; import com.opensource.docgrid.domain.user.entity.User; import com.opensource.docgrid.domain.user.repository.UserRepository; import com.opensource.docgrid.global.exception.DocGridException; @@ -40,6 +42,7 @@ public class DocumentUploadService { private final DocumentVersionRepository documentVersionRepository; private final EmbeddingJobRepository embeddingJobRepository; private final EmbeddingModelQueryService embeddingModelQueryService; + private final SyncEventWriter syncEventWriter; @Transactional(readOnly = true) public Optional findReusableFileObjectId(String fileHash, long fileSize) { @@ -84,10 +87,13 @@ public DocumentUploadTransactionResult upload(DocumentUploadCommand command) { document.updateCurrentVersion(documentVersion); EmbeddingModel embeddingModel = embeddingModelQueryService.getActiveModel(); + // Version과 최초 인덱싱 의도를 같은 Transaction에 기록해 둘 중 하나만 Commit되는 상태를 막는다. + SyncOutboxEvent sourceEvent = syncEventWriter.recordDocumentVersionCreated(documentVersion, embeddingModel); EmbeddingJob embeddingJob = embeddingJobRepository.save( EmbeddingJob.builder() .documentVersion(documentVersion) .embeddingModel(embeddingModel) + .sourceEventId(sourceEvent.getEventId()) .status(EmbeddingJobStatus.PENDING) .priority(DEFAULT_JOB_PRIORITY) .maxRetryCount(MAX_RETRY_COUNT) diff --git a/backend/src/main/java/com/opensource/docgrid/domain/document/service/command/DocumentVersionUploadService.java b/backend/src/main/java/com/opensource/docgrid/domain/document/service/command/DocumentVersionUploadService.java index b082c8c..ed52487 100644 --- a/backend/src/main/java/com/opensource/docgrid/domain/document/service/command/DocumentVersionUploadService.java +++ b/backend/src/main/java/com/opensource/docgrid/domain/document/service/command/DocumentVersionUploadService.java @@ -22,6 +22,8 @@ import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus; import com.opensource.docgrid.domain.embedding.repository.EmbeddingJobRepository; import com.opensource.docgrid.domain.embedding.service.query.EmbeddingModelQueryService; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.service.command.SyncEventWriter; import com.opensource.docgrid.domain.user.entity.User; import com.opensource.docgrid.domain.user.repository.UserRepository; import com.opensource.docgrid.global.exception.DocGridException; @@ -50,6 +52,7 @@ public class DocumentVersionUploadService { private final DocumentVersionRepository documentVersionRepository; private final EmbeddingJobRepository embeddingJobRepository; private final EmbeddingModelQueryService embeddingModelQueryService; + private final SyncEventWriter syncEventWriter; @Transactional(readOnly = true) public Optional prepare(Long userId, Long documentId, ValidatedFile file, String fileHash) { @@ -100,10 +103,13 @@ public DocumentVersionUploadTransactionResult upload(DocumentVersionUploadComman } EmbeddingModel embeddingModel = embeddingModelQueryService.getActiveModel(); + // 새 Version과 인덱싱 의도를 같은 Transaction에 저장해 후속 복구의 기준 Event를 남긴다. + SyncOutboxEvent sourceEvent = syncEventWriter.recordDocumentVersionCreated(version, embeddingModel); EmbeddingJob job = embeddingJobRepository.save( EmbeddingJob.builder() .documentVersion(version) .embeddingModel(embeddingModel) + .sourceEventId(sourceEvent.getEventId()) .status(EmbeddingJobStatus.PENDING) .priority(DEFAULT_JOB_PRIORITY) .maxRetryCount(MAX_RETRY_COUNT) diff --git a/backend/src/main/java/com/opensource/docgrid/domain/embedding/entity/EmbeddingJob.java b/backend/src/main/java/com/opensource/docgrid/domain/embedding/entity/EmbeddingJob.java index f95022a..d2a94a1 100644 --- a/backend/src/main/java/com/opensource/docgrid/domain/embedding/entity/EmbeddingJob.java +++ b/backend/src/main/java/com/opensource/docgrid/domain/embedding/entity/EmbeddingJob.java @@ -1,6 +1,7 @@ package com.opensource.docgrid.domain.embedding.entity; import java.time.LocalDateTime; +import java.util.UUID; import com.opensource.docgrid.domain.document.entity.DocumentVersion; import com.opensource.docgrid.domain.embedding.enums.EmbeddingJobStatus; @@ -33,6 +34,7 @@ * -> Worker가 파싱/청킹/임베딩 -> embeddings 저장). * 관계: document_version_id -> DocumentVersion, embedding_model_id -> EmbeddingModel, * locked_by_worker_id -> WorkerNode(nullable, lock을 잡은 Worker). + * source_event_id는 이 Job을 처음 만든 Sync Outbox Event를 가리키며 Event 재전달의 중복 Job 생성을 막는다. * index: (status, priority, created_at) 우선순위 큐 조회용, (status, next_retry_at) Retry 실행 가능 시각 조회용, * lock_expires_at, (document_version_id, embedding_model_id), locked_by_worker_id. * @@ -118,15 +120,20 @@ public class EmbeddingJob extends BaseEntity { @Column(name = "error_message", columnDefinition = "TEXT") private String errorMessage; + // 최초 Job 생성 원인을 식별하며 UNIQUE 제약으로 같은 Event의 중복 Job 생성을 차단한다. + @Column(name = "source_event_id", unique = true) + private UUID sourceEventId; + @Builder public EmbeddingJob(DocumentVersion documentVersion, EmbeddingModel embeddingModel, EmbeddingJobStatus status, - int priority, int maxRetryCount) { + int priority, int maxRetryCount, UUID sourceEventId) { this.documentVersion = documentVersion; this.embeddingModel = embeddingModel; this.status = status != null ? status : EmbeddingJobStatus.PENDING; this.priority = priority; this.retryCount = 0; this.maxRetryCount = maxRetryCount; + this.sourceEventId = sourceEventId; } /** 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 new file mode 100644 index 0000000..bb5fed9 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/entity/SyncOutboxEvent.java @@ -0,0 +1,144 @@ +package com.opensource.docgrid.domain.sync.entity; + +import java.time.LocalDateTime; +import java.util.UUID; + +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.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; + +/** + * 도메인 변경과 같은 Transaction에서 생성되는 동기화 Outbox Event다. + * + *

역할: 문서·권한·모델 변경을 후속 Dispatcher가 유실 없이 발견할 수 있는 영속 Queue로 보존한다. + * eventId는 전달 세대를, idempotencyKey는 같은 비즈니스 변경의 중복 생성을 식별한다. Payload는 처리 + * 힌트일 뿐이며 Handler는 aggregateType과 aggregateId로 현재 DB 상태를 다시 읽어야 한다. + * + *

경계: 인덱싱 상태를 사람이 조회하는 append-only 이력은 {@code IndexingEvent}가 담당하고, 이 + * Entity는 실제로 소비·재시도해야 하는 작업만 관리한다. Dispatcher Claim과 상태 전이는 후속 단계에서 + * 이 Entity의 Lease 필드를 이용한다. + */ +@Getter +@Entity +@NoArgsConstructor(access = AccessLevel.PROTECTED) +@Table( + name = "sync_outbox_events", + uniqueConstraints = { + @UniqueConstraint(name = "uk_sync_outbox_events_event_id", columnNames = "event_id"), + @UniqueConstraint(name = "uk_sync_outbox_events_idempotency_key", columnNames = "idempotency_key") + }, + indexes = { + @Index(name = "idx_sync_outbox_events_dispatch", columnList = "status, available_at, occurred_at, id"), + @Index(name = "idx_sync_outbox_events_lock_expires_at", columnList = "lock_expires_at"), + @Index( + name = "idx_sync_outbox_events_aggregate", + columnList = "aggregate_type, aggregate_id, aggregate_version" + ) + } +) +public class SyncOutboxEvent extends BaseEntity { + + @Id + @GeneratedValue(strategy = GenerationType.IDENTITY) + private Long id; + + @Column(name = "event_id", nullable = false, updatable = false) + private UUID eventId; + + @Column(name = "idempotency_key", nullable = false, updatable = false, length = 255) + private String idempotencyKey; + + @Enumerated(EnumType.STRING) + @Column(name = "aggregate_type", nullable = false, updatable = false, length = 30) + private SyncAggregateType aggregateType; + + @Column(name = "aggregate_id", nullable = false, updatable = false) + private Long aggregateId; + + @Column(name = "aggregate_version", updatable = false) + private Long aggregateVersion; + + @Enumerated(EnumType.STRING) + @Column(name = "event_type", nullable = false, updatable = false, length = 50) + private SyncEventType eventType; + + // OpenSQL과 현재 Hibernate 설정의 기존 JSON 경계를 따라 JSON 문자열을 TEXT로 보존한다. + @Column(name = "payload_json", columnDefinition = "TEXT", updatable = false) + private String payloadJson; + + @Enumerated(EnumType.STRING) + @Column(nullable = false, length = 20) + private SyncEventStatus status; + + @Column(name = "available_at", nullable = false) + private LocalDateTime availableAt; + + @Column(name = "occurred_at", nullable = false, updatable = false) + private LocalDateTime occurredAt; + + @Column(name = "processed_at") + private LocalDateTime processedAt; + + @Column(name = "retry_count", nullable = false) + private int retryCount; + + @Column(name = "max_retry_count", nullable = false) + private int maxRetryCount; + + @Column(name = "claim_token") + private UUID claimToken; + + @Column(name = "locked_by", length = 200) + private String lockedBy; + + @Column(name = "lock_expires_at") + private LocalDateTime lockExpiresAt; + + @Column(name = "last_error_code", length = 100) + private String lastErrorCode; + + @Column(name = "last_error_message", columnDefinition = "TEXT") + private String lastErrorMessage; + + @Builder + public SyncOutboxEvent( + UUID eventId, + String idempotencyKey, + SyncAggregateType aggregateType, + Long aggregateId, + Long aggregateVersion, + SyncEventType eventType, + String payloadJson, + LocalDateTime availableAt, + LocalDateTime occurredAt, + int maxRetryCount + ) { + this.eventId = eventId; + this.idempotencyKey = idempotencyKey; + this.aggregateType = aggregateType; + this.aggregateId = aggregateId; + this.aggregateVersion = aggregateVersion; + this.eventType = eventType; + this.payloadJson = payloadJson; + this.status = SyncEventStatus.PENDING; + this.availableAt = availableAt; + this.occurredAt = occurredAt; + this.maxRetryCount = maxRetryCount; + } +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncAggregateType.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncAggregateType.java new file mode 100644 index 0000000..3bbd376 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncAggregateType.java @@ -0,0 +1,14 @@ +package com.opensource.docgrid.domain.sync.enums; + +/** + * 동기화 Event가 설명하는 도메인 Aggregate의 경계를 정의한다. + * + *

Dispatcher는 이 값과 aggregateId를 조합해 현재 도메인 상태를 다시 읽으며, Event Payload를 + * 최신 상태의 원장으로 사용하지 않는다. + */ +public enum SyncAggregateType { + DOCUMENT_VERSION, + DOCUMENT, + PERMISSION, + EMBEDDING_MODEL +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncEventStatus.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncEventStatus.java new file mode 100644 index 0000000..9a1a2df --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncEventStatus.java @@ -0,0 +1,14 @@ +package com.opensource.docgrid.domain.sync.enums; + +/** + * Outbox Event의 전달 생명주기를 표현한다. + * + *

PENDING과 PROCESSING은 Dispatcher의 Lease 소유권 경계이고, PROCESSED와 FAILED는 후속 처리가 + * 끝난 종결 상태다. + */ +public enum SyncEventStatus { + PENDING, + PROCESSING, + PROCESSED, + FAILED +} diff --git a/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncEventType.java b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncEventType.java new file mode 100644 index 0000000..69b01d7 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/enums/SyncEventType.java @@ -0,0 +1,15 @@ +package com.opensource.docgrid.domain.sync.enums; + +/** + * Outbox에서 전달할 후속 동기화 작업의 종류를 정의한다. + * + *

각 값은 Dispatcher Handler 하나의 멱등한 책임에 대응하며, 인덱싱 운영 이력인 + * {@code indexing_events}와 분리된다. + */ +public enum SyncEventType { + DOCUMENT_VERSION_CREATED, + DOCUMENT_REINDEX_REQUESTED, + DOCUMENT_DELETED, + PERMISSION_CACHE_REFRESH_REQUESTED, + EMBEDDING_MODEL_ACTIVATED +} 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 new file mode 100644 index 0000000..50b98cf --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/repository/SyncOutboxEventRepository.java @@ -0,0 +1,20 @@ +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.SyncOutboxEvent; + +/** + * 동기화 Outbox Event의 영속성과 멱등 식별자 조회를 담당한다. + * + *

Dispatcher의 Claim·Lease Query는 해당 동작을 구현하는 단계에서 이 Repository에 추가한다. + */ +public interface SyncOutboxEventRepository extends JpaRepository { + + Optional findByEventId(UUID eventId); + + Optional findByIdempotencyKey(String idempotencyKey); +} 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 new file mode 100644 index 0000000..e5c57c2 --- /dev/null +++ b/backend/src/main/java/com/opensource/docgrid/domain/sync/service/command/SyncEventWriter.java @@ -0,0 +1,68 @@ +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.document.entity.DocumentVersion; +import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel; +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.repository.SyncOutboxEventRepository; + +import lombok.RequiredArgsConstructor; + +/** + * 도메인 변경 Transaction 안에서 후속 동기화 Outbox Event를 기록한다. + * + *

이 Service는 Event를 외부로 발행하지 않는다. 호출한 Command Transaction과 같은 경계에서 Event를 + * 저장해 도메인 변경만 Commit되거나 Event만 남는 이중 쓰기 상태를 방지하는 것이 책임이다. + */ +@Service +@RequiredArgsConstructor +@Transactional +public class SyncEventWriter { + + private static final int DEFAULT_MAX_RETRY_COUNT = 5; + + private final SyncOutboxEventRepository syncOutboxEventRepository; + private final Clock clock; + + /** + * 새 문서 버전의 최초 인덱싱 의도를 같은 Transaction에 기록한다. + */ + public SyncOutboxEvent recordDocumentVersionCreated( + DocumentVersion documentVersion, + EmbeddingModel embeddingModel + ) { + LocalDateTime occurredAt = LocalDateTime.now(clock); + String idempotencyKey = String.format( + "%s:%d:%s:%d", + SyncAggregateType.DOCUMENT_VERSION, + documentVersion.getId(), + SyncEventType.DOCUMENT_VERSION_CREATED, + documentVersion.getVersionNo() + ); + String payloadJson = String.format("{\"embeddingModelId\":%d}", embeddingModel.getId()); + + // Version·Job 생성 Transaction과 함께 Commit돼야 Dispatcher가 부분 상태를 관측하지 않는다. + return syncOutboxEventRepository.save( + SyncOutboxEvent.builder() + .eventId(UUID.randomUUID()) + .idempotencyKey(idempotencyKey) + .aggregateType(SyncAggregateType.DOCUMENT_VERSION) + .aggregateId(documentVersion.getId()) + .aggregateVersion((long) documentVersion.getVersionNo()) + .eventType(SyncEventType.DOCUMENT_VERSION_CREATED) + .payloadJson(payloadJson) + .occurredAt(occurredAt) + .availableAt(occurredAt) + .maxRetryCount(DEFAULT_MAX_RETRY_COUNT) + .build() + ); + } +} diff --git a/backend/src/main/resources/db/migration/V36__create_sync_outbox_events.sql b/backend/src/main/resources/db/migration/V36__create_sync_outbox_events.sql new file mode 100644 index 0000000..f31459b --- /dev/null +++ b/backend/src/main/resources/db/migration/V36__create_sync_outbox_events.sql @@ -0,0 +1,41 @@ +-- sync_outbox_events: 도메인 변경을 후속 동기화 처리로 안전하게 전달하는 Transactional Outbox +CREATE TABLE sync_outbox_events ( + id BIGSERIAL PRIMARY KEY, + event_id UUID NOT NULL, + idempotency_key VARCHAR(255) NOT NULL, + aggregate_type VARCHAR(30) NOT NULL, + aggregate_id BIGINT NOT NULL, + aggregate_version BIGINT, + event_type VARCHAR(50) NOT NULL, + payload_json TEXT, + status VARCHAR(20) NOT NULL, + available_at TIMESTAMP NOT NULL, + occurred_at TIMESTAMP NOT NULL, + processed_at TIMESTAMP, + retry_count INT NOT NULL DEFAULT 0, + max_retry_count INT NOT NULL, + claim_token UUID, + locked_by VARCHAR(200), + lock_expires_at TIMESTAMP, + last_error_code VARCHAR(100), + last_error_message TEXT, + created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + + CONSTRAINT uk_sync_outbox_events_event_id UNIQUE (event_id), + CONSTRAINT uk_sync_outbox_events_idempotency_key UNIQUE (idempotency_key) +); + +CREATE INDEX idx_sync_outbox_events_dispatch + ON sync_outbox_events (status, available_at, occurred_at, id); +CREATE INDEX idx_sync_outbox_events_lock_expires_at + ON sync_outbox_events (lock_expires_at); +CREATE INDEX idx_sync_outbox_events_aggregate + ON sync_outbox_events (aggregate_type, aggregate_id, aggregate_version); + +-- 최초 Job과 원인이 된 Event를 연결해 Event 재전달에서도 Job을 한 번만 만들 수 있게 한다. +ALTER TABLE embedding_jobs + ADD COLUMN source_event_id UUID, + ADD CONSTRAINT fk_embedding_jobs_source_event + FOREIGN KEY (source_event_id) REFERENCES sync_outbox_events (event_id), + ADD CONSTRAINT uk_embedding_jobs_source_event_id UNIQUE (source_event_id); diff --git a/backend/src/test/java/com/opensource/docgrid/domain/document/integration/DocumentUploadIntegrationTest.java b/backend/src/test/java/com/opensource/docgrid/domain/document/integration/DocumentUploadIntegrationTest.java index 0e99ffc..ae1ab68 100644 --- a/backend/src/test/java/com/opensource/docgrid/domain/document/integration/DocumentUploadIntegrationTest.java +++ b/backend/src/test/java/com/opensource/docgrid/domain/document/integration/DocumentUploadIntegrationTest.java @@ -79,6 +79,10 @@ void setUp() { void tearDown() { for (DocumentUploadResponse response : createdResponses) { jdbcTemplate.update("DELETE FROM embedding_jobs WHERE id = ?", response.embeddingJobId()); + jdbcTemplate.update( + "DELETE FROM sync_outbox_events WHERE aggregate_type = 'DOCUMENT_VERSION' AND aggregate_id = ?", + response.documentVersionId() + ); jdbcTemplate.update("UPDATE documents SET current_version_id = NULL WHERE id = ?", response.documentId()); jdbcTemplate.update("DELETE FROM document_versions WHERE id = ?", response.documentVersionId()); jdbcTemplate.update("DELETE FROM documents WHERE id = ?", response.documentId()); @@ -118,6 +122,20 @@ void upload_createsAllReceptionData() { assertThat(jdbcTemplate.queryForObject( "SELECT max_retry_count FROM embedding_jobs WHERE id = ?", Integer.class, response.embeddingJobId() )).isEqualTo(3); + UUID sourceEventId = jdbcTemplate.queryForObject( + "SELECT source_event_id FROM embedding_jobs WHERE id = ?", + UUID.class, + response.embeddingJobId() + ); + assertThat(sourceEventId).isNotNull(); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM sync_outbox_events " + + "WHERE event_id = ? AND aggregate_type = 'DOCUMENT_VERSION' AND aggregate_id = ? " + + "AND event_type = 'DOCUMENT_VERSION_CREATED' AND status = 'PENDING'", + Integer.class, + sourceEventId, + response.documentVersionId() + )).isOne(); } @Test diff --git a/backend/src/test/java/com/opensource/docgrid/domain/document/integration/DocumentVersionUploadIntegrationTest.java b/backend/src/test/java/com/opensource/docgrid/domain/document/integration/DocumentVersionUploadIntegrationTest.java index 7b60422..74844dd 100644 --- a/backend/src/test/java/com/opensource/docgrid/domain/document/integration/DocumentVersionUploadIntegrationTest.java +++ b/backend/src/test/java/com/opensource/docgrid/domain/document/integration/DocumentVersionUploadIntegrationTest.java @@ -78,6 +78,11 @@ void tearDown() { ); jdbcTemplate.update("DELETE FROM embedding_jobs WHERE document_version_id IN " + "(SELECT id FROM document_versions WHERE document_id = ?)", documentId); + jdbcTemplate.update( + "DELETE FROM sync_outbox_events WHERE aggregate_type = 'DOCUMENT_VERSION' " + + "AND aggregate_id IN (SELECT id FROM document_versions WHERE document_id = ?)", + documentId + ); jdbcTemplate.update("UPDATE documents SET current_version_id = NULL WHERE id = ?", documentId); jdbcTemplate.update("DELETE FROM document_versions WHERE document_id = ?", documentId); jdbcTemplate.update("DELETE FROM documents WHERE id = ?", documentId); @@ -118,6 +123,13 @@ void upload_createsVersion_when_sameSizeButContentDiffers() { String.class, response.documentVersionId() )).isEqualTo("text/plain"); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM sync_outbox_events " + + "WHERE aggregate_type = 'DOCUMENT_VERSION' AND aggregate_id = ? " + + "AND event_type = 'DOCUMENT_VERSION_CREATED' AND status = 'PENDING'", + Integer.class, + response.documentVersionId() + )).isOne(); } @Test diff --git a/backend/src/test/java/com/opensource/docgrid/domain/document/service/command/DocumentUploadServiceTest.java b/backend/src/test/java/com/opensource/docgrid/domain/document/service/command/DocumentUploadServiceTest.java index 3333fa3..986d733 100644 --- a/backend/src/test/java/com/opensource/docgrid/domain/document/service/command/DocumentUploadServiceTest.java +++ b/backend/src/test/java/com/opensource/docgrid/domain/document/service/command/DocumentUploadServiceTest.java @@ -7,6 +7,7 @@ import static org.mockito.BDDMockito.then; import java.util.Optional; +import java.util.UUID; import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.DisplayName; @@ -35,6 +36,8 @@ import com.opensource.docgrid.domain.embedding.fixture.EmbeddingModelFixture; import com.opensource.docgrid.domain.embedding.repository.EmbeddingJobRepository; import com.opensource.docgrid.domain.embedding.service.query.EmbeddingModelQueryService; +import com.opensource.docgrid.domain.sync.entity.SyncOutboxEvent; +import com.opensource.docgrid.domain.sync.service.command.SyncEventWriter; import com.opensource.docgrid.domain.user.entity.User; import com.opensource.docgrid.domain.user.enums.UserStatus; import com.opensource.docgrid.domain.user.repository.UserRepository; @@ -47,6 +50,7 @@ class DocumentUploadServiceTest { private static final Long USER_ID = 1L; private static final String FILE_HASH = "hash"; + private static final UUID SOURCE_EVENT_ID = UUID.fromString("00000000-0000-0000-0000-000000000162"); @Mock private UserRepository userRepository; @Mock private FileObjectRepository fileObjectRepository; @@ -54,6 +58,8 @@ class DocumentUploadServiceTest { @Mock private DocumentVersionRepository documentVersionRepository; @Mock private EmbeddingJobRepository embeddingJobRepository; @Mock private EmbeddingModelQueryService embeddingModelQueryService; + @Mock private SyncEventWriter syncEventWriter; + @Mock private SyncOutboxEvent syncOutboxEvent; private DocumentUploadService documentUploadService; private User user; @@ -68,7 +74,8 @@ void setUp() { documentRepository, documentVersionRepository, embeddingJobRepository, - embeddingModelQueryService + embeddingModelQueryService, + syncEventWriter ); user = User.builder() .email("user@test.com") @@ -102,6 +109,7 @@ void upload_createsDocumentVersionAndJob_withExistingFileObject() { given(documentVersionRepository.save(any(DocumentVersion.class))) .willAnswer(invocation -> withId(invocation.getArgument(0), 11L)); given(embeddingModelQueryService.getActiveModel()).willReturn(embeddingModel); + givenSourceEvent(); given(embeddingJobRepository.save(any(EmbeddingJob.class))) .willAnswer(invocation -> withId(invocation.getArgument(0), 12L)); @@ -130,6 +138,8 @@ void upload_createsDocumentVersionAndJob_withExistingFileObject() { assertThat(jobCaptor.getValue().getRetryCount()).isZero(); assertThat(jobCaptor.getValue().getMaxRetryCount()).isEqualTo(3); assertThat(jobCaptor.getValue().getEmbeddingModel()).isSameAs(embeddingModel); + assertThat(jobCaptor.getValue().getSourceEventId()).isEqualTo(SOURCE_EVENT_ID); + then(syncEventWriter).should().recordDocumentVersionCreated(versionCaptor.getValue(), embeddingModel); } @Test @@ -145,6 +155,7 @@ void upload_reusesFileObject_when_atomicInsertLosesRace() { given(documentVersionRepository.save(any(DocumentVersion.class))) .willAnswer(invocation -> withId(invocation.getArgument(0), 11L)); given(embeddingModelQueryService.getActiveModel()).willReturn(embeddingModel); + givenSourceEvent(); given(embeddingJobRepository.save(any(EmbeddingJob.class))) .willAnswer(invocation -> withId(invocation.getArgument(0), 12L)); @@ -181,6 +192,12 @@ private DocumentUploadCommand command(Long existingFileObjectId, StoredFile stor ); } + private void givenSourceEvent() { + given(syncEventWriter.recordDocumentVersionCreated(any(DocumentVersion.class), any(EmbeddingModel.class))) + .willReturn(syncOutboxEvent); + given(syncOutboxEvent.getEventId()).willReturn(SOURCE_EVENT_ID); + } + private T withId(T entity, Long id) { ReflectionTestUtils.setField(entity, "id", id); return entity; diff --git a/backend/src/test/java/com/opensource/docgrid/domain/sync/service/command/SyncEventWriterTest.java b/backend/src/test/java/com/opensource/docgrid/domain/sync/service/command/SyncEventWriterTest.java new file mode 100644 index 0000000..4be3f73 --- /dev/null +++ b/backend/src/test/java/com/opensource/docgrid/domain/sync/service/command/SyncEventWriterTest.java @@ -0,0 +1,86 @@ +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 java.time.Clock; +import java.time.Instant; +import java.time.LocalDateTime; +import java.time.ZoneId; + +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.DocumentVersion; +import com.opensource.docgrid.domain.document.enums.DocumentVersionStatus; +import com.opensource.docgrid.domain.embedding.entity.EmbeddingModel; +import com.opensource.docgrid.domain.embedding.fixture.EmbeddingModelFixture; +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; + +/** + * SyncEventWriter가 문서 버전 변경을 재현 가능한 멱등 Outbox Event로 변환하는지 검증한다. + */ +@ExtendWith(MockitoExtension.class) +@DisplayName("SyncEventWriter 단위 테스트") +class SyncEventWriterTest { + + private static final Instant FIXED_INSTANT = Instant.parse("2026-08-13T06:00:00Z"); + private static final ZoneId ZONE_ID = ZoneId.of("Asia/Seoul"); + + @Mock private SyncOutboxEventRepository syncOutboxEventRepository; + + private SyncEventWriter syncEventWriter; + + @BeforeEach + void setUp() { + syncEventWriter = new SyncEventWriter( + syncOutboxEventRepository, + Clock.fixed(FIXED_INSTANT, ZONE_ID) + ); + } + + @Test + @DisplayName("문서 버전과 모델을 PENDING Outbox Event로 기록한다") + void recordDocumentVersionCreated_savesPendingEvent() { + DocumentVersion version = DocumentVersion.builder() + .versionNo(3) + .status(DocumentVersionStatus.UPLOADED) + .build(); + ReflectionTestUtils.setField(version, "id", 341L); + EmbeddingModel model = EmbeddingModelFixture.createDefaultModel(); + ReflectionTestUtils.setField(model, "id", 7L); + given(syncOutboxEventRepository.save(any(SyncOutboxEvent.class))) + .willAnswer(invocation -> invocation.getArgument(0)); + + SyncOutboxEvent result = syncEventWriter.recordDocumentVersionCreated(version, model); + + ArgumentCaptor eventCaptor = ArgumentCaptor.forClass(SyncOutboxEvent.class); + then(syncOutboxEventRepository).should().save(eventCaptor.capture()); + assertThat(result).isSameAs(eventCaptor.getValue()); + assertThat(result.getEventId()).isNotNull(); + assertThat(result.getIdempotencyKey()) + .isEqualTo("DOCUMENT_VERSION:341:DOCUMENT_VERSION_CREATED:3"); + assertThat(result.getAggregateType()).isEqualTo(SyncAggregateType.DOCUMENT_VERSION); + assertThat(result.getAggregateId()).isEqualTo(341L); + assertThat(result.getAggregateVersion()).isEqualTo(3L); + assertThat(result.getEventType()).isEqualTo(SyncEventType.DOCUMENT_VERSION_CREATED); + assertThat(result.getPayloadJson()).isEqualTo("{\"embeddingModelId\":7}"); + assertThat(result.getStatus()).isEqualTo(SyncEventStatus.PENDING); + assertThat(result.getAvailableAt()).isEqualTo(LocalDateTime.ofInstant(FIXED_INSTANT, ZONE_ID)); + assertThat(result.getOccurredAt()).isEqualTo(LocalDateTime.ofInstant(FIXED_INSTANT, ZONE_ID)); + assertThat(result.getRetryCount()).isZero(); + assertThat(result.getMaxRetryCount()).isEqualTo(5); + } +}