Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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<Long> findReusableFileObjectId(String fileHash, long fileSize) {
Expand Down Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<Long> prepare(Long userId, Long documentId, ValidatedFile file, String fileHash) {
Expand Down Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
@@ -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;
Expand Down Expand Up @@ -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.
*
Expand Down Expand Up @@ -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;
}

/**
Expand Down
Original file line number Diff line number Diff line change
@@ -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다.
*
* <p>역할: 문서·권한·모델 변경을 후속 Dispatcher가 유실 없이 발견할 수 있는 영속 Queue로 보존한다.
* eventId는 전달 세대를, idempotencyKey는 같은 비즈니스 변경의 중복 생성을 식별한다. Payload는 처리
* 힌트일 뿐이며 Handler는 aggregateType과 aggregateId로 현재 DB 상태를 다시 읽어야 한다.
*
* <p>경계: 인덱싱 상태를 사람이 조회하는 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;
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
package com.opensource.docgrid.domain.sync.enums;

/**
* 동기화 Event가 설명하는 도메인 Aggregate의 경계를 정의한다.
*
* <p>Dispatcher는 이 값과 aggregateId를 조합해 현재 도메인 상태를 다시 읽으며, Event Payload를
* 최신 상태의 원장으로 사용하지 않는다.
*/
public enum SyncAggregateType {
DOCUMENT_VERSION,
DOCUMENT,
PERMISSION,
EMBEDDING_MODEL
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
package com.opensource.docgrid.domain.sync.enums;

/**
* Outbox Event의 전달 생명주기를 표현한다.
*
* <p>PENDING과 PROCESSING은 Dispatcher의 Lease 소유권 경계이고, PROCESSED와 FAILED는 후속 처리가
* 끝난 종결 상태다.
*/
public enum SyncEventStatus {
PENDING,
PROCESSING,
PROCESSED,
FAILED
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
package com.opensource.docgrid.domain.sync.enums;

/**
* Outbox에서 전달할 후속 동기화 작업의 종류를 정의한다.
*
* <p>각 값은 Dispatcher Handler 하나의 멱등한 책임에 대응하며, 인덱싱 운영 이력인
* {@code indexing_events}와 분리된다.
*/
public enum SyncEventType {
DOCUMENT_VERSION_CREATED,
DOCUMENT_REINDEX_REQUESTED,
DOCUMENT_DELETED,
PERMISSION_CACHE_REFRESH_REQUESTED,
EMBEDDING_MODEL_ACTIVATED
}
Original file line number Diff line number Diff line change
@@ -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의 영속성과 멱등 식별자 조회를 담당한다.
*
* <p>Dispatcher의 Claim·Lease Query는 해당 동작을 구현하는 단계에서 이 Repository에 추가한다.
*/
public interface SyncOutboxEventRepository extends JpaRepository<SyncOutboxEvent, Long> {

Optional<SyncOutboxEvent> findByEventId(UUID eventId);

Optional<SyncOutboxEvent> findByIdempotencyKey(String idempotencyKey);
}
Original file line number Diff line number Diff line change
@@ -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를 기록한다.
*
* <p>이 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()
);
}
}
Loading