diff --git a/docs/reliability/transactional-outbox.md b/docs/reliability/transactional-outbox.md index 6ff48cc..1cf02b0 100644 --- a/docs/reliability/transactional-outbox.md +++ b/docs/reliability/transactional-outbox.md @@ -24,6 +24,7 @@ Kafka나 RabbitMQ가 필요한 구조는 아닙니다. MVP는 기존 PostgreSQL | --- | --- | | `event_publication` | 처리해야 할 이벤트와 현재 상태, 시도 횟수, 다음 시각, lease를 저장 | | `event_consumption` | 어느 handler가 어느 이벤트를 이미 성공했는지 저장 | +| `outbox_manual_retry` | ADMIN의 수동 재처리 사유와 중복 방지 키 해시를 변경 불가 기록으로 저장 | `event_publication`의 상태는 다음과 같습니다. @@ -128,8 +129,36 @@ ORDER BY occurred_at; `PROCESSING` lease가 만료되면 다음 poll에서 자동 복구됩니다. `RETRY_WAIT`도 `next_attempt_at` 이후 자동 처리됩니다. `REVIEW_REQUIRED`는 원인을 수정했다고 해서 -DB를 임의로 `PENDING`으로 바꾸지 않습니다. 별도 관리 command와 감사로그가 구현되기 -전에는 담당 개발자가 원인을 확인하고 forward migration 또는 후속 Issue로 복구합니다. +DB를 임의로 `PENDING`으로 바꾸지 않습니다. + +### REVIEW_REQUIRED 수동 재처리 + +1. 안전한 `last_error_code`와 관련 handler 상태를 확인하고 원인을 먼저 해결합니다. +2. ADMIN Access Token으로 아래 API를 호출합니다. +3. 응답이 `202`이면 이벤트는 `PENDING`이 되고 다음 Outbox poll에서 한 번 더 시도됩니다. +4. `event_publication`, `outbox_manual_retry`, `audit_event`를 payload 없이 확인합니다. + +```http +POST /api/v1/admin/outbox-events/{eventId}/retry +Authorization: Bearer {admin-access-token} +Idempotency-Key: incident-20260806-event-001 +Content-Type: application/json + +{ + "expected_version": 3, + "reason": "내부 handler 복구와 점검을 완료했습니다." +} +``` + +- `expected_version`은 조회 당시 `event_publication.version`입니다. 그 사이 상태가 바뀌면 + `409`로 거부하므로 최신 상태를 다시 확인합니다. +- 같은 `Idempotency-Key`와 같은 요청은 재실행하지 않고 최초 접수 결과를 반환합니다. +- 사유는 10~300자로 작성하며 이름·연락처·문서 원문·token·payload·예외 원문을 + 입력하지 않습니다. +- 수동 재처리가 승인되면 `attempt_count`를 0으로 초기화해 handler가 실제로 다시 + 실행될 기회를 부여합니다. 초기화 전 횟수는 `outbox_manual_retry.previous_attempt_count`에 + 변경 불가 이력으로 보존합니다. 이후 실패하면 일반 자동 재시도 한도를 다시 적용합니다. +- HR·VIEWER와 다른 사업장의 ADMIN은 호출할 수 없습니다. `OUTBOX_ENABLED=false`는 자동 처리를 멈출 뿐 새 이벤트 저장을 막지 않습니다. 장애 중 이벤트가 계속 누적될 수 있으므로 backlog를 함께 관찰하고, 수정 배포 후 다시 활성화해 @@ -143,6 +172,8 @@ DB를 임의로 `PENDING`으로 바꾸지 않습니다. 별도 관리 command와 constraint, index - 기능 통합 테스트: 실제 command가 올바른 event type과 최소 payload를 발행하는지 검증 +- 운영 API 통합 테스트: ADMIN 권한, 사업장 격리, version 충돌, Idempotency-Key, + 동시 재처리와 감사로그, 최대 횟수 이벤트의 실제 handler 재실행 검증 로컬 전체 검증: diff --git a/src/main/java/com/fowoco/server/audit/domain/AuditAction.java b/src/main/java/com/fowoco/server/audit/domain/AuditAction.java index 2aa5117..7e433d4 100644 --- a/src/main/java/com/fowoco/server/audit/domain/AuditAction.java +++ b/src/main/java/com/fowoco/server/audit/domain/AuditAction.java @@ -18,6 +18,7 @@ public enum AuditAction { AI_RUN_CREATED, AI_RUN_ANSWERS_SUBMITTED, AI_RUN_CANDIDATES_DECIDED, + OUTBOX_MANUAL_RETRY_REQUESTED, WORKER_LINK_RESPONSE_SUBMITTED, WORKER_LINK_ACCESSED, USER_AGREEMENTS_RECORDED, diff --git a/src/main/java/com/fowoco/server/audit/domain/AuditTargetType.java b/src/main/java/com/fowoco/server/audit/domain/AuditTargetType.java index eb84c86..291cd5c 100644 --- a/src/main/java/com/fowoco/server/audit/domain/AuditTargetType.java +++ b/src/main/java/com/fowoco/server/audit/domain/AuditTargetType.java @@ -9,5 +9,6 @@ public enum AuditTargetType { WORKER_DOCUMENT, DOCUMENT_REQUEST_DRAFT, AI_RUN, + OUTBOX_EVENT, USER_ACCOUNT } diff --git a/src/main/java/com/fowoco/server/common/error/GlobalExceptionHandler.java b/src/main/java/com/fowoco/server/common/error/GlobalExceptionHandler.java index 3dee917..4e5cf62 100644 --- a/src/main/java/com/fowoco/server/common/error/GlobalExceptionHandler.java +++ b/src/main/java/com/fowoco/server/common/error/GlobalExceptionHandler.java @@ -26,6 +26,7 @@ import org.springframework.web.HttpMediaTypeNotSupportedException; import org.springframework.web.HttpRequestMethodNotSupportedException; import org.springframework.web.bind.MethodArgumentNotValidException; +import org.springframework.web.bind.MissingRequestHeaderException; import org.springframework.web.bind.MissingServletRequestParameterException; import org.springframework.web.bind.annotation.ExceptionHandler; import org.springframework.web.bind.annotation.RestControllerAdvice; @@ -104,6 +105,7 @@ public ResponseEntity handleMethodValidation( @ExceptionHandler({ HttpMessageNotReadableException.class, + MissingRequestHeaderException.class, MissingServletRequestParameterException.class, MethodArgumentTypeMismatchException.class }) diff --git a/src/main/java/com/fowoco/server/reliability/api/OutboxManualRetryRequest.java b/src/main/java/com/fowoco/server/reliability/api/OutboxManualRetryRequest.java new file mode 100644 index 0000000..3af5d58 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/api/OutboxManualRetryRequest.java @@ -0,0 +1,25 @@ +package com.fowoco.server.reliability.api; + +import com.fasterxml.jackson.annotation.JsonProperty; +import io.swagger.v3.oas.annotations.media.Schema; +import jakarta.validation.constraints.Min; +import jakarta.validation.constraints.NotBlank; +import jakarta.validation.constraints.NotNull; +import jakarta.validation.constraints.Size; + +public record OutboxManualRetryRequest( + @JsonProperty("expected_version") + @Schema(description = "운영자가 확인한 현재 Outbox event version", example = "3") + @NotNull(message = "expected_version은 필수입니다.") + @Min(value = 0, message = "expected_version은 0 이상이어야 합니다.") + Long expectedVersion, + + @Schema( + description = "재처리 근거. 개인정보·payload·token·예외 원문은 입력하지 않습니다.", + example = "일시 중단된 내부 handler 복구를 확인함" + ) + @NotBlank(message = "재처리 사유를 입력해 주세요.") + @Size(min = 10, max = 300, message = "재처리 사유는 10자 이상 300자 이하로 입력해 주세요.") + String reason +) { +} diff --git a/src/main/java/com/fowoco/server/reliability/api/OutboxManualRetryResponse.java b/src/main/java/com/fowoco/server/reliability/api/OutboxManualRetryResponse.java new file mode 100644 index 0000000..8c5ae59 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/api/OutboxManualRetryResponse.java @@ -0,0 +1,25 @@ +package com.fowoco.server.reliability.api; + +import com.fasterxml.jackson.annotation.JsonProperty; +import com.fowoco.server.reliability.application.OutboxManualRetryResult; +import com.fowoco.server.reliability.domain.EventPublicationStatus; +import java.time.Instant; +import java.util.UUID; + +public record OutboxManualRetryResponse( + @JsonProperty("event_id") UUID eventId, + @JsonProperty("accepted_status") EventPublicationStatus acceptedStatus, + @JsonProperty("accepted_version") long acceptedVersion, + @JsonProperty("accepted_at") Instant acceptedAt, + @JsonProperty("already_requested") boolean alreadyRequested +) { + static OutboxManualRetryResponse from(OutboxManualRetryResult result) { + return new OutboxManualRetryResponse( + result.eventId(), + result.acceptedStatus(), + result.acceptedVersion(), + result.acceptedAt(), + result.alreadyRequested() + ); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/api/OutboxOperationsController.java b/src/main/java/com/fowoco/server/reliability/api/OutboxOperationsController.java new file mode 100644 index 0000000..ac68803 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/api/OutboxOperationsController.java @@ -0,0 +1,81 @@ +package com.fowoco.server.reliability.api; + +import com.fowoco.server.auth.application.port.ActorContextProvider; +import com.fowoco.server.common.web.RequestMetadata; +import com.fowoco.server.reliability.application.OutboxManualRetryService; +import io.swagger.v3.oas.annotations.Operation; +import io.swagger.v3.oas.annotations.Parameter; +import io.swagger.v3.oas.annotations.responses.ApiResponse; +import io.swagger.v3.oas.annotations.responses.ApiResponses; +import io.swagger.v3.oas.annotations.security.SecurityRequirement; +import io.swagger.v3.oas.annotations.tags.Tag; +import jakarta.servlet.http.HttpServletRequest; +import jakarta.validation.Valid; +import java.util.UUID; +import org.springframework.http.MediaType; +import org.springframework.http.ResponseEntity; +import org.springframework.security.access.prepost.PreAuthorize; +import org.springframework.web.bind.annotation.PathVariable; +import org.springframework.web.bind.annotation.PostMapping; +import org.springframework.web.bind.annotation.RequestBody; +import org.springframework.web.bind.annotation.RequestHeader; +import org.springframework.web.bind.annotation.RequestMapping; +import org.springframework.web.bind.annotation.RestController; + +@Tag(name = "Operations", description = "ADMIN 전용 장애 복구 작업") +@SecurityRequirement(name = "bearerAuth") +@RestController +@RequestMapping("/api/v1/admin/outbox-events") +public class OutboxOperationsController { + + private final OutboxManualRetryService retryService; + private final ActorContextProvider actorContextProvider; + + public OutboxOperationsController( + OutboxManualRetryService retryService, + ActorContextProvider actorContextProvider + ) { + this.retryService = retryService; + this.actorContextProvider = actorContextProvider; + } + + @Operation( + operationId = "retryOutboxEvent", + summary = "REVIEW_REQUIRED Outbox 이벤트 재처리", + description = "ADMIN이 원인을 해결한 뒤 멈춘 이벤트를 PENDING으로 되돌립니다. " + + "payload·예외 원문은 반환하지 않으며 실제 처리는 Outbox worker가 수행합니다." + ) + @ApiResponses({ + @ApiResponse(responseCode = "202", description = "재처리 요청 접수"), + @ApiResponse(responseCode = "400", ref = "#/components/responses/BadRequest"), + @ApiResponse(responseCode = "401", ref = "#/components/responses/Unauthorized"), + @ApiResponse(responseCode = "403", ref = "#/components/responses/Forbidden"), + @ApiResponse(responseCode = "404", ref = "#/components/responses/NotFound"), + @ApiResponse(responseCode = "409", ref = "#/components/responses/Conflict") + }) + @PreAuthorize("hasRole('ADMIN')") + @PostMapping( + path = "/{eventId}/retry", + consumes = MediaType.APPLICATION_JSON_VALUE, + produces = MediaType.APPLICATION_JSON_VALUE + ) + public ResponseEntity retry( + @Parameter(description = "재처리할 Outbox event ID") + @PathVariable UUID eventId, + @Parameter(description = "같은 운영 요청의 중복 실행을 막는 키", required = true) + @RequestHeader("Idempotency-Key") String idempotencyKey, + @Valid @RequestBody OutboxManualRetryRequest request, + HttpServletRequest servletRequest + ) { + return ResponseEntity.accepted().body(OutboxManualRetryResponse.from( + retryService.requestRetry( + eventId, + request.expectedVersion(), + request.reason(), + idempotencyKey, + actorContextProvider.requireCurrentActor(), + RequestMetadata.from(servletRequest) + ) + )); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxManualRetryResult.java b/src/main/java/com/fowoco/server/reliability/application/OutboxManualRetryResult.java new file mode 100644 index 0000000..d350d6a --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/OutboxManualRetryResult.java @@ -0,0 +1,14 @@ +package com.fowoco.server.reliability.application; + +import com.fowoco.server.reliability.domain.EventPublicationStatus; +import java.time.Instant; +import java.util.UUID; + +public record OutboxManualRetryResult( + UUID eventId, + EventPublicationStatus acceptedStatus, + long acceptedVersion, + Instant acceptedAt, + boolean alreadyRequested +) { +} diff --git a/src/main/java/com/fowoco/server/reliability/application/OutboxManualRetryService.java b/src/main/java/com/fowoco/server/reliability/application/OutboxManualRetryService.java new file mode 100644 index 0000000..8092811 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/OutboxManualRetryService.java @@ -0,0 +1,194 @@ +package com.fowoco.server.reliability.application; + +import com.fowoco.server.audit.application.port.AuditEventRepository; +import com.fowoco.server.audit.domain.ActorType; +import com.fowoco.server.audit.domain.AuditAction; +import com.fowoco.server.audit.domain.AuditEvent; +import com.fowoco.server.audit.domain.AuditTargetType; +import com.fowoco.server.auth.application.ActorAuthorizer; +import com.fowoco.server.auth.application.ActorContext; +import com.fowoco.server.auth.domain.UserRole; +import com.fowoco.server.common.error.ApiException; +import com.fowoco.server.common.error.ErrorCode; +import com.fowoco.server.common.id.UuidGenerator; +import com.fowoco.server.common.security.TenantDatabaseContext; +import com.fowoco.server.common.web.RequestMetadata; +import com.fowoco.server.reliability.application.error.OutboxErrorCode; +import com.fowoco.server.reliability.application.port.EventPublicationRepository; +import com.fowoco.server.reliability.application.port.OutboxManualRetryRepository; +import com.fowoco.server.reliability.domain.EventPublication; +import com.fowoco.server.reliability.domain.EventPublicationStatus; +import com.fowoco.server.reliability.domain.OutboxManualRetry; +import java.nio.charset.StandardCharsets; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; +import java.time.Clock; +import java.time.Instant; +import java.util.UUID; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +@Service +public class OutboxManualRetryService { + + private static final String AUDIT_EVENT_VERSION = "1.0"; + + private final ActorAuthorizer actorAuthorizer; + private final TenantDatabaseContext tenantDatabaseContext; + private final EventPublicationRepository publicationRepository; + private final OutboxManualRetryRepository retryRepository; + private final AuditEventRepository auditEventRepository; + private final UuidGenerator uuidGenerator; + private final Clock clock; + + public OutboxManualRetryService( + ActorAuthorizer actorAuthorizer, + TenantDatabaseContext tenantDatabaseContext, + EventPublicationRepository publicationRepository, + OutboxManualRetryRepository retryRepository, + AuditEventRepository auditEventRepository, + UuidGenerator uuidGenerator, + Clock clock + ) { + this.actorAuthorizer = actorAuthorizer; + this.tenantDatabaseContext = tenantDatabaseContext; + this.publicationRepository = publicationRepository; + this.retryRepository = retryRepository; + this.auditEventRepository = auditEventRepository; + this.uuidGenerator = uuidGenerator; + this.clock = clock; + } + + @Transactional + public OutboxManualRetryResult requestRetry( + UUID eventId, + long expectedVersion, + String reason, + String idempotencyKey, + ActorContext actor, + RequestMetadata metadata + ) { + actorAuthorizer.requireAnyRole(actor, UserRole.ADMIN); + tenantDatabaseContext.setCompanyIdForCurrentTransaction(actor.companyId()); + + String normalizedReason = normalizeReason(reason); + String keyHash = sha256(normalizeIdempotencyKey(idempotencyKey)); + String requestHash = sha256(expectedVersion + "\n" + normalizedReason); + OutboxManualRetry existing = retryRepository + .findByEventIdAndCompanyIdAndKeyHash(eventId, actor.companyId(), keyHash) + .orElse(null); + if (existing != null) { + return replay(existing, requestHash); + } + + EventPublication publication = publicationRepository + .findByIdAndCompanyIdForUpdate(eventId, actor.companyId()) + .orElseThrow(() -> new ApiException(OutboxErrorCode.OUTBOX_EVENT_NOT_FOUND)); + + existing = retryRepository + .findByEventIdAndCompanyIdAndKeyHash(eventId, actor.companyId(), keyHash) + .orElse(null); + if (existing != null) { + return replay(existing, requestHash); + } + if (publication.version() != expectedVersion) { + throw new ApiException(OutboxErrorCode.OUTBOX_EVENT_VERSION_CONFLICT); + } + if (publication.status() != EventPublicationStatus.REVIEW_REQUIRED + || publication.leaseOwner() != null + || publication.leaseExpiresAt() != null) { + throw new ApiException(OutboxErrorCode.OUTBOX_EVENT_NOT_REVIEW_REQUIRED); + } + + Instant now = clock.instant(); + int previousAttemptCount = publication.attemptCount(); + publication.requestManualRetry(expectedVersion, now); + EventPublication saved = publicationRepository.save(publication); + OutboxManualRetry retry = retryRepository.append(new OutboxManualRetry( + uuidGenerator.generate(), + actor.companyId(), + eventId, + keyHash, + requestHash, + normalizedReason, + actor.actorId(), + metadata.requestId(), + metadata.traceId(), + previousAttemptCount, + saved.status(), + saved.version(), + now + )); + appendAudit(retry, actor, metadata, now); + + return result(retry, false); + } + + private OutboxManualRetryResult replay(OutboxManualRetry existing, String requestHash) { + if (!existing.requestHash().equals(requestHash)) { + throw new ApiException(OutboxErrorCode.OUTBOX_RETRY_IDEMPOTENCY_CONFLICT); + } + return result(existing, true); + } + + private OutboxManualRetryResult result(OutboxManualRetry retry, boolean alreadyRequested) { + return new OutboxManualRetryResult( + retry.eventId(), + retry.acceptedStatus(), + retry.acceptedVersion(), + retry.createdAt(), + alreadyRequested + ); + } + + private String normalizeReason(String reason) { + if (reason == null) { + throw new ApiException(ErrorCode.VALIDATION_FAILED); + } + String normalized = reason.trim(); + if (normalized.length() < 10 || normalized.length() > 300) { + throw new ApiException(ErrorCode.VALIDATION_FAILED); + } + return normalized; + } + + private String normalizeIdempotencyKey(String key) { + if (key == null || key.isBlank() || key.length() > 100) { + throw new ApiException(OutboxErrorCode.OUTBOX_RETRY_INVALID_IDEMPOTENCY_KEY); + } + return key.trim(); + } + + private String sha256(String value) { + try { + byte[] digest = MessageDigest.getInstance("SHA-256") + .digest(value.getBytes(StandardCharsets.UTF_8)); + return java.util.HexFormat.of().formatHex(digest); + } catch (NoSuchAlgorithmException exception) { + throw new IllegalStateException("SHA-256 must be available", exception); + } + } + + private void appendAudit( + OutboxManualRetry retry, + ActorContext actor, + RequestMetadata metadata, + Instant now + ) { + auditEventRepository.append(new AuditEvent( + uuidGenerator.generate(), + actor.companyId(), + ActorType.HR_USER, + actor.actorId(), + UserRole.ADMIN, + AuditAction.OUTBOX_MANUAL_RETRY_REQUESTED, + AuditTargetType.OUTBOX_EVENT, + retry.eventId(), + metadata.requestId(), + metadata.traceId(), + AUDIT_EVENT_VERSION, + "REVIEW_REQUIRED 이벤트 재처리 요청: " + retry.reason(), + now + )); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/error/OutboxErrorCode.java b/src/main/java/com/fowoco/server/reliability/application/error/OutboxErrorCode.java new file mode 100644 index 0000000..e2418e9 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/error/OutboxErrorCode.java @@ -0,0 +1,47 @@ +package com.fowoco.server.reliability.application.error; + +import com.fowoco.server.common.error.ApiErrorCode; +import org.springframework.http.HttpStatus; + +public enum OutboxErrorCode implements ApiErrorCode { + OUTBOX_EVENT_NOT_FOUND(HttpStatus.NOT_FOUND, "Outbox 이벤트를 찾을 수 없습니다."), + OUTBOX_EVENT_NOT_REVIEW_REQUIRED( + HttpStatus.CONFLICT, + "수동 확인이 필요한 Outbox 이벤트만 다시 처리할 수 있습니다." + ), + OUTBOX_EVENT_VERSION_CONFLICT( + HttpStatus.CONFLICT, + "이벤트 상태가 먼저 변경되었습니다. 최신 정보를 다시 확인해 주세요." + ), + OUTBOX_RETRY_IDEMPOTENCY_CONFLICT( + HttpStatus.CONFLICT, + "같은 Idempotency-Key가 다른 재처리 요청에 이미 사용되었습니다." + ), + OUTBOX_RETRY_INVALID_IDEMPOTENCY_KEY( + HttpStatus.BAD_REQUEST, + "Idempotency-Key를 확인해 주세요." + ); + + private final HttpStatus status; + private final String defaultMessage; + + OutboxErrorCode(HttpStatus status, String defaultMessage) { + this.status = status; + this.defaultMessage = defaultMessage; + } + + @Override + public String code() { + return name(); + } + + @Override + public HttpStatus status() { + return status; + } + + @Override + public String defaultMessage() { + return defaultMessage; + } +} diff --git a/src/main/java/com/fowoco/server/reliability/application/port/OutboxManualRetryRepository.java b/src/main/java/com/fowoco/server/reliability/application/port/OutboxManualRetryRepository.java new file mode 100644 index 0000000..d2c85dd --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/application/port/OutboxManualRetryRepository.java @@ -0,0 +1,16 @@ +package com.fowoco.server.reliability.application.port; + +import com.fowoco.server.reliability.domain.OutboxManualRetry; +import java.util.Optional; +import java.util.UUID; + +public interface OutboxManualRetryRepository { + + Optional findByEventIdAndCompanyIdAndKeyHash( + UUID eventId, + UUID companyId, + String idempotencyKeyHash + ); + + OutboxManualRetry append(OutboxManualRetry retry); +} diff --git a/src/main/java/com/fowoco/server/reliability/domain/EventPublication.java b/src/main/java/com/fowoco/server/reliability/domain/EventPublication.java index 599a0d2..c1eb3e1 100644 --- a/src/main/java/com/fowoco/server/reliability/domain/EventPublication.java +++ b/src/main/java/com/fowoco/server/reliability/domain/EventPublication.java @@ -228,6 +228,24 @@ public void requireReview(String owner, String errorCode, Instant now) { updatedAt = now; } + public void requestManualRetry(long expectedVersion, Instant now) { + Objects.requireNonNull(now); + if (version != expectedVersion) { + throw new IllegalArgumentException("Event publication version does not match."); + } + if (status != EventPublicationStatus.REVIEW_REQUIRED + || leaseOwner != null + || leaseExpiresAt != null) { + throw new IllegalStateException("Only an unleased review-required event can be retried."); + } + status = EventPublicationStatus.PENDING; + attemptCount = 0; + nextAttemptAt = now; + lastErrorCode = null; + completedAt = null; + updatedAt = now; + } + public void requireActiveLease(String owner, Instant now) { String normalizedOwner = requireOwner(owner); Objects.requireNonNull(now); diff --git a/src/main/java/com/fowoco/server/reliability/domain/OutboxManualRetry.java b/src/main/java/com/fowoco/server/reliability/domain/OutboxManualRetry.java new file mode 100644 index 0000000..d461f56 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/domain/OutboxManualRetry.java @@ -0,0 +1,56 @@ +package com.fowoco.server.reliability.domain; + +import java.time.Instant; +import java.util.Objects; +import java.util.UUID; + +public record OutboxManualRetry( + UUID manualRetryId, + UUID companyId, + UUID eventId, + String idempotencyKeyHash, + String requestHash, + String reason, + UUID requestedBy, + String requestId, + String traceId, + int previousAttemptCount, + EventPublicationStatus acceptedStatus, + long acceptedVersion, + Instant createdAt +) { + public OutboxManualRetry { + Objects.requireNonNull(manualRetryId, "manualRetryId must not be null"); + Objects.requireNonNull(companyId, "companyId must not be null"); + Objects.requireNonNull(eventId, "eventId must not be null"); + Objects.requireNonNull(requestedBy, "requestedBy must not be null"); + Objects.requireNonNull(acceptedStatus, "acceptedStatus must not be null"); + Objects.requireNonNull(createdAt, "createdAt must not be null"); + requireLength(idempotencyKeyHash, 64, "idempotencyKeyHash"); + requireLength(requestHash, 64, "requestHash"); + requireRange(reason, 10, 300, "reason"); + requireRange(requestId, 1, 128, "requestId"); + if (traceId != null) { + requireLength(traceId, 32, "traceId"); + } + if (acceptedVersion < 0) { + throw new IllegalArgumentException("acceptedVersion must not be negative"); + } + if (previousAttemptCount < 0) { + throw new IllegalArgumentException("previousAttemptCount must not be negative"); + } + } + + private static void requireLength(String value, int length, String name) { + if (value == null || value.length() != length) { + throw new IllegalArgumentException(name + " length is invalid"); + } + } + + private static void requireRange(String value, int minimum, int maximum, String name) { + if (value == null || value.isBlank() + || value.length() < minimum || value.length() > maximum) { + throw new IllegalArgumentException(name + " length is invalid"); + } + } +} diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaOutboxManualRetryRepository.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaOutboxManualRetryRepository.java new file mode 100644 index 0000000..3166441 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/JpaOutboxManualRetryRepository.java @@ -0,0 +1,36 @@ +package com.fowoco.server.reliability.infrastructure.persistence; + +import com.fowoco.server.reliability.application.port.OutboxManualRetryRepository; +import com.fowoco.server.reliability.domain.OutboxManualRetry; +import java.util.Optional; +import java.util.UUID; +import org.springframework.stereotype.Repository; + +@Repository +public class JpaOutboxManualRetryRepository implements OutboxManualRetryRepository { + + private final SpringDataOutboxManualRetryJpaRepository repository; + + public JpaOutboxManualRetryRepository(SpringDataOutboxManualRetryJpaRepository repository) { + this.repository = repository; + } + + @Override + public Optional findByEventIdAndCompanyIdAndKeyHash( + UUID eventId, + UUID companyId, + String idempotencyKeyHash + ) { + return repository.findByEventIdAndCompanyIdAndIdempotencyKeyHash( + eventId, + companyId, + idempotencyKeyHash + ) + .map(OutboxManualRetryJpaEntity::toDomain); + } + + @Override + public OutboxManualRetry append(OutboxManualRetry retry) { + return repository.saveAndFlush(new OutboxManualRetryJpaEntity(retry)).toDomain(); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/OutboxManualRetryJpaEntity.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/OutboxManualRetryJpaEntity.java new file mode 100644 index 0000000..4ce7517 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/OutboxManualRetryJpaEntity.java @@ -0,0 +1,83 @@ +package com.fowoco.server.reliability.infrastructure.persistence; + +import com.fowoco.server.reliability.domain.EventPublicationStatus; +import com.fowoco.server.reliability.domain.OutboxManualRetry; +import jakarta.persistence.Column; +import jakarta.persistence.Entity; +import jakarta.persistence.EnumType; +import jakarta.persistence.Enumerated; +import jakarta.persistence.Id; +import jakarta.persistence.Table; +import java.time.Instant; +import java.util.UUID; + +@Entity +@Table(name = "outbox_manual_retry") +class OutboxManualRetryJpaEntity { + + @Id + @Column(name = "manual_retry_id", nullable = false, updatable = false) + private UUID manualRetryId; + @Column(name = "company_id", nullable = false, updatable = false) + private UUID companyId; + @Column(name = "event_id", nullable = false, updatable = false) + private UUID eventId; + @Column(name = "idempotency_key_hash", nullable = false, length = 64, updatable = false) + private String idempotencyKeyHash; + @Column(name = "request_hash", nullable = false, length = 64, updatable = false) + private String requestHash; + @Column(name = "reason", nullable = false, length = 300, updatable = false) + private String reason; + @Column(name = "requested_by", nullable = false, updatable = false) + private UUID requestedBy; + @Column(name = "request_id", nullable = false, length = 128, updatable = false) + private String requestId; + @Column(name = "trace_id", length = 32, updatable = false) + private String traceId; + @Column(name = "previous_attempt_count", nullable = false, updatable = false) + private int previousAttemptCount; + @Enumerated(EnumType.STRING) + @Column(name = "accepted_status", nullable = false, length = 30, updatable = false) + private EventPublicationStatus acceptedStatus; + @Column(name = "accepted_version", nullable = false, updatable = false) + private long acceptedVersion; + @Column(name = "created_at", nullable = false, updatable = false) + private Instant createdAt; + + protected OutboxManualRetryJpaEntity() { + } + + OutboxManualRetryJpaEntity(OutboxManualRetry retry) { + manualRetryId = retry.manualRetryId(); + companyId = retry.companyId(); + eventId = retry.eventId(); + idempotencyKeyHash = retry.idempotencyKeyHash(); + requestHash = retry.requestHash(); + reason = retry.reason(); + requestedBy = retry.requestedBy(); + requestId = retry.requestId(); + traceId = retry.traceId(); + previousAttemptCount = retry.previousAttemptCount(); + acceptedStatus = retry.acceptedStatus(); + acceptedVersion = retry.acceptedVersion(); + createdAt = retry.createdAt(); + } + + OutboxManualRetry toDomain() { + return new OutboxManualRetry( + manualRetryId, + companyId, + eventId, + idempotencyKeyHash, + requestHash, + reason, + requestedBy, + requestId, + traceId, + previousAttemptCount, + acceptedStatus, + acceptedVersion, + createdAt + ); + } +} diff --git a/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/SpringDataOutboxManualRetryJpaRepository.java b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/SpringDataOutboxManualRetryJpaRepository.java new file mode 100644 index 0000000..5e59480 --- /dev/null +++ b/src/main/java/com/fowoco/server/reliability/infrastructure/persistence/SpringDataOutboxManualRetryJpaRepository.java @@ -0,0 +1,15 @@ +package com.fowoco.server.reliability.infrastructure.persistence; + +import java.util.Optional; +import java.util.UUID; +import org.springframework.data.jpa.repository.JpaRepository; + +interface SpringDataOutboxManualRetryJpaRepository + extends JpaRepository { + + Optional findByEventIdAndCompanyIdAndIdempotencyKeyHash( + UUID eventId, + UUID companyId, + String idempotencyKeyHash + ); +} diff --git a/src/main/resources/db/migration-postgresql/V27__prepare_outbox_manual_retry_rls.sql b/src/main/resources/db/migration-postgresql/V27__prepare_outbox_manual_retry_rls.sql new file mode 100644 index 0000000..685a252 --- /dev/null +++ b/src/main/resources/db/migration-postgresql/V27__prepare_outbox_manual_retry_rls.sql @@ -0,0 +1,12 @@ +CREATE POLICY pl_outbox_manual_retry_tenant_isolation + ON public.outbox_manual_retry + FOR ALL + TO PUBLIC + USING ( + company_id = + NULLIF(pg_catalog.current_setting('app.company_id', true), '')::UUID + ) + WITH CHECK ( + company_id = + NULLIF(pg_catalog.current_setting('app.company_id', true), '')::UUID + ); diff --git a/src/main/resources/db/migration/V26__create_outbox_manual_retry.sql b/src/main/resources/db/migration/V26__create_outbox_manual_retry.sql new file mode 100644 index 0000000..815d4b0 --- /dev/null +++ b/src/main/resources/db/migration/V26__create_outbox_manual_retry.sql @@ -0,0 +1,46 @@ +CREATE TABLE outbox_manual_retry ( + manual_retry_id UUID NOT NULL, + company_id UUID NOT NULL, + event_id UUID NOT NULL, + idempotency_key_hash VARCHAR(64) NOT NULL, + request_hash VARCHAR(64) NOT NULL, + reason VARCHAR(300) NOT NULL, + requested_by UUID NOT NULL, + request_id VARCHAR(128) NOT NULL, + trace_id VARCHAR(32), + previous_attempt_count INTEGER NOT NULL DEFAULT 0, + accepted_status VARCHAR(30) NOT NULL, + accepted_version BIGINT NOT NULL, + created_at TIMESTAMP(6) WITH TIME ZONE NOT NULL, + CONSTRAINT pk_outbox_manual_retry PRIMARY KEY (manual_retry_id), + CONSTRAINT uq_outbox_manual_retry_event_key + UNIQUE (company_id, event_id, idempotency_key_hash), + CONSTRAINT fk_outbox_manual_retry_event_company + FOREIGN KEY (event_id, company_id) + REFERENCES event_publication (event_id, company_id) ON DELETE RESTRICT, + CONSTRAINT fk_outbox_manual_retry_actor_company + FOREIGN KEY (requested_by, company_id) + REFERENCES user_account (user_id, company_id) ON DELETE RESTRICT, + CONSTRAINT ck_outbox_manual_retry_key_hash + CHECK (CHAR_LENGTH(idempotency_key_hash) = 64), + CONSTRAINT ck_outbox_manual_retry_request_hash + CHECK (CHAR_LENGTH(request_hash) = 64), + CONSTRAINT ck_outbox_manual_retry_reason + CHECK (CHAR_LENGTH(TRIM(reason)) BETWEEN 10 AND 300), + CONSTRAINT ck_outbox_manual_retry_request_id + CHECK (CHAR_LENGTH(TRIM(request_id)) BETWEEN 1 AND 128), + CONSTRAINT ck_outbox_manual_retry_trace_id + CHECK (trace_id IS NULL OR CHAR_LENGTH(trace_id) = 32), + CONSTRAINT ck_outbox_manual_retry_previous_attempt_count + CHECK (previous_attempt_count >= 0), + CONSTRAINT ck_outbox_manual_retry_status + CHECK (accepted_status = 'PENDING'), + CONSTRAINT ck_outbox_manual_retry_version + CHECK (accepted_version >= 0) +); + +CREATE INDEX idx_outbox_manual_retry_company_created + ON outbox_manual_retry (company_id, created_at); + +CREATE INDEX idx_outbox_manual_retry_event_created + ON outbox_manual_retry (company_id, event_id, created_at); diff --git a/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java b/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java index 69eb4bc..0dfbf55 100644 --- a/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java +++ b/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java @@ -88,6 +88,7 @@ private void assertSchemaContract(Connection connection) throws SQLException { "audit_event", "event_publication", "event_consumption", + "outbox_manual_retry", "document_request_draft", "document_request_draft_type", "ai_run", @@ -204,6 +205,16 @@ private void assertSchemaContract(Connection connection) throws SQLException { .containsEntry("company_id", new ColumnSpec("uuid", false)) .containsEntry("handler_name", new ColumnSpec("varchar", false)) .containsEntry("completed_at", new ColumnSpec("timestamptz", false)); + assertThat(columnSpecs(connection, "outbox_manual_retry")) + .containsEntry("manual_retry_id", new ColumnSpec("uuid", false)) + .containsEntry("company_id", new ColumnSpec("uuid", false)) + .containsEntry("event_id", new ColumnSpec("uuid", false)) + .containsEntry("idempotency_key_hash", new ColumnSpec("varchar", false)) + .containsEntry("request_hash", new ColumnSpec("varchar", false)) + .containsEntry("reason", new ColumnSpec("varchar", false)) + .containsEntry("requested_by", new ColumnSpec("uuid", false)) + .containsEntry("previous_attempt_count", new ColumnSpec("int4", false)) + .containsEntry("accepted_version", new ColumnSpec("int8", false)); assertThat(columnSpecs(connection, "document_request_draft")) .containsEntry("draft_id", new ColumnSpec("uuid", false)) .containsEntry("task_id", new ColumnSpec("uuid", false)) @@ -342,6 +353,10 @@ private void assertSchemaContract(Connection connection) throws SQLException { "fk_worker_response_upload_file_company", "fk_worker_document_upload_idempotency_link_company", "fk_worker_document_upload_idempotency_file_company", + "pk_outbox_manual_retry", + "uq_outbox_manual_retry_event_key", + "fk_outbox_manual_retry_event_company", + "fk_outbox_manual_retry_actor_company", "pk_user_agreement_consent", "fk_user_agreement_consent_user_company", "pk_password_reset_token", @@ -376,6 +391,8 @@ private void assertSchemaContract(Connection connection) throws SQLException { "idx_worker_response_upload_company", "idx_worker_document_upload_idempotency_company", "idx_worker_document_upload_idempotency_file_company", + "idx_outbox_manual_retry_company_created", + "idx_outbox_manual_retry_event_created", "idx_user_agreement_consent_user_time", "idx_password_reset_token_company_user", "idx_password_reset_token_active" @@ -397,6 +414,7 @@ private void assertSchemaContract(Connection connection) throws SQLException { "pl_audit_event_tenant_isolation", "pl_event_publication_tenant_isolation", "pl_event_consumption_tenant_isolation", + "pl_outbox_manual_retry_tenant_isolation", "pl_document_request_draft_tenant_isolation", "pl_document_request_draft_type_tenant_isolation", "pl_ai_run_tenant_isolation", diff --git a/src/test/java/com/fowoco/server/common/security/PostgreSqlRlsIsolationTest.java b/src/test/java/com/fowoco/server/common/security/PostgreSqlRlsIsolationTest.java index 4ab6a2a..ffc4a2e 100644 --- a/src/test/java/com/fowoco/server/common/security/PostgreSqlRlsIsolationTest.java +++ b/src/test/java/com/fowoco/server/common/security/PostgreSqlRlsIsolationTest.java @@ -62,6 +62,14 @@ class PostgreSqlRlsIsolationTest { UUID.fromString("b7000000-0000-0000-0000-000000000002"); private static final UUID CASE_A_NEW = UUID.fromString("a7000000-0000-0000-0000-000000000003"); + private static final UUID EVENT_A = + UUID.fromString("a9000000-0000-0000-0000-000000000001"); + private static final UUID EVENT_B = + UUID.fromString("b9000000-0000-0000-0000-000000000002"); + private static final UUID MANUAL_RETRY_A = + UUID.fromString("aa000000-0000-0000-0000-000000000001"); + private static final UUID MANUAL_RETRY_B = + UUID.fromString("bb000000-0000-0000-0000-000000000002"); private static final UUID CONSENT_A = UUID.fromString("a9000000-0000-0000-0000-000000000001"); private static final UUID CONSENT_B = @@ -83,6 +91,7 @@ class PostgreSqlRlsIsolationTest { "worker_response", "worker_response_upload", "worker_document_upload_idempotency", + "outbox_manual_retry", "user_agreement_consent", "password_reset_token" ); @@ -176,6 +185,7 @@ private void prepareFixture( + "public.worker_link, public.worker_response, " + "public.worker_response_upload, " + "public.worker_document_upload_idempotency, " + + "public.outbox_manual_retry, " + "public.user_agreement_consent, " + "public.password_reset_token TO " + quotedRole @@ -339,6 +349,41 @@ INSERT INTO worker_document_upload_idempotency ( WORKER_LINK_A, COMPANY_A, STORED_FILE_A, WORKER_LINK_B, COMPANY_B, STORED_FILE_B )); + statement.execute(""" + INSERT INTO event_publication ( + event_id, company_id, event_type, payload_version, + aggregate_type, aggregate_id, actor_type, request_id, + payload_json, status, attempt_count, last_error_code, + occurred_at, created_at, updated_at, version + ) VALUES + ('%s', '%s', 'RlsEvent', '1', 'RlsProbe', '%s', + 'SYSTEM_RULE', 'rls-event-a', '{}', 'REVIEW_REQUIRED', 3, + 'RLS_REVIEW_REQUIRED', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, + CURRENT_TIMESTAMP, 0), + ('%s', '%s', 'RlsEvent', '1', 'RlsProbe', '%s', + 'SYSTEM_RULE', 'rls-event-b', '{}', 'REVIEW_REQUIRED', 3, + 'RLS_REVIEW_REQUIRED', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, + CURRENT_TIMESTAMP, 0) + """.formatted( + EVENT_A, COMPANY_A, UUID.randomUUID(), + EVENT_B, COMPANY_B, UUID.randomUUID() + )); + statement.execute(""" + INSERT INTO outbox_manual_retry ( + manual_retry_id, company_id, event_id, + idempotency_key_hash, request_hash, reason, requested_by, + request_id, accepted_status, accepted_version, created_at + ) VALUES + ('%s', '%s', '%s', repeat('a', 64), repeat('c', 64), + 'Tenant A handler 복구 확인', '%s', 'rls-retry-a', + 'PENDING', 1, CURRENT_TIMESTAMP), + ('%s', '%s', '%s', repeat('b', 64), repeat('d', 64), + 'Tenant B handler 복구 확인', '%s', 'rls-retry-b', + 'PENDING', 1, CURRENT_TIMESTAMP) + """.formatted( + MANUAL_RETRY_A, COMPANY_A, EVENT_A, USER_A, + MANUAL_RETRY_B, COMPANY_B, EVENT_B, USER_B + )); } rlsState.enableRowLevelSecurity(); @@ -357,6 +402,7 @@ private void assertMissingAndInvalidContextFailClosed(Connection connection) assertThat(tableCount(connection, "worker_response_upload")).isZero(); assertThat(tableCount(connection, "worker_document_upload_idempotency")).isZero(); assertThat(tableCount(connection, "workflow_case")).isZero(); + assertThat(tableCount(connection, "outbox_manual_retry")).isZero(); assertThat(tableCount(connection, "user_agreement_consent")).isZero(); assertThat(tableCount(connection, "password_reset_token")).isZero(); @@ -417,6 +463,11 @@ private void assertTenantCrudIsolation(Connection connection) throws SQLExceptio connection, "SELECT case_id FROM public.workflow_case ORDER BY case_id" )).containsExactly(CASE_A); + assertThat(uuidValues( + connection, + "SELECT manual_retry_id FROM public.outbox_manual_retry " + + "ORDER BY manual_retry_id" + )).containsExactly(MANUAL_RETRY_A); assertThat(uuidValues( connection, "SELECT consent_id FROM public.user_agreement_consent ORDER BY consent_id" @@ -567,6 +618,22 @@ INSERT INTO worker_response_upload ( STORED_FILE_B_UNLINKED, COMPANY_B )); + assertSqlState( + connection, + "42501", + """ + INSERT INTO outbox_manual_retry ( + manual_retry_id, company_id, event_id, + idempotency_key_hash, request_hash, reason, requested_by, + request_id, accepted_status, accepted_version, created_at + ) VALUES ( + 'bb000000-0000-0000-0000-000000000099', '%s', '%s', + repeat('e', 64), repeat('f', 64), + 'Forbidden tenant retry request', '%s', 'rls-forbidden-retry', + 'PENDING', 1, CURRENT_TIMESTAMP + ) + """.formatted(COMPANY_B, EVENT_B, USER_B) + ); assertThat(executeUpdate( connection, @@ -593,6 +660,11 @@ INSERT INTO worker_response_upload ( "DELETE FROM worker WHERE worker_id = ?", WORKER_B )).isZero(); + assertThat(executeUpdate( + connection, + "DELETE FROM outbox_manual_retry WHERE manual_retry_id = ?", + MANUAL_RETRY_B + )).isZero(); assertThat(executeUpdate( connection, "DELETE FROM worker WHERE worker_id = ?", @@ -617,12 +689,17 @@ private void assertCommittedContextDoesNotLeak(Connection connection) throws SQL assertThat(workerCount(connection)).isZero(); assertThat(tableCount(connection, "workflow_case")).isZero(); + assertThat(tableCount(connection, "outbox_manual_retry")).isZero(); setTenantContext(connection, COMPANY_B.toString()); assertThat(workerIds(connection)).containsExactly(WORKER_B); assertThat(uuidValues( connection, "SELECT case_id FROM public.workflow_case ORDER BY case_id" )).containsExactly(CASE_B); + assertThat(uuidValues( + connection, + "SELECT manual_retry_id FROM public.outbox_manual_retry ORDER BY manual_retry_id" + )).containsExactly(MANUAL_RETRY_B); connection.rollback(); } @@ -665,6 +742,14 @@ private SQLException runCleanupStep(SQLException failure, SqlCleanupStep step) { } private void deleteFixtureRows(Statement statement) throws SQLException { + statement.execute(""" + DELETE FROM outbox_manual_retry + WHERE manual_retry_id IN ('%s', '%s') + """.formatted(MANUAL_RETRY_A, MANUAL_RETRY_B)); + statement.execute(""" + DELETE FROM event_publication + WHERE event_id IN ('%s', '%s') + """.formatted(EVENT_A, EVENT_B)); statement.execute(""" DELETE FROM password_reset_token WHERE password_reset_token_id IN ('%s', '%s') diff --git a/src/test/java/com/fowoco/server/reliability/OutboxManualRetryApiIntegrationTest.java b/src/test/java/com/fowoco/server/reliability/OutboxManualRetryApiIntegrationTest.java new file mode 100644 index 0000000..53c9052 --- /dev/null +++ b/src/test/java/com/fowoco/server/reliability/OutboxManualRetryApiIntegrationTest.java @@ -0,0 +1,487 @@ +package com.fowoco.server.reliability; + +import static org.assertj.core.api.Assertions.assertThat; + +import com.fowoco.server.reliability.application.OutboxProcessor; +import com.fowoco.server.reliability.application.port.DomainEventHandler; +import com.fowoco.server.reliability.domain.DomainEventEnvelope; +import com.jayway.jsonpath.JsonPath; +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.HttpRequest; +import java.net.http.HttpResponse; +import java.util.List; +import java.util.UUID; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.boot.test.context.TestConfiguration; +import org.springframework.boot.test.web.server.LocalServerPort; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Import; +import org.springframework.http.HttpHeaders; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.security.crypto.password.PasswordEncoder; +import org.springframework.test.context.ActiveProfiles; + +@ActiveProfiles("test") +@Import(OutboxManualRetryApiIntegrationTest.ManualRetryTestConfiguration.class) +@SpringBootTest( + webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT, + properties = "app.reliability.outbox.enabled=false" +) +class OutboxManualRetryApiIntegrationTest { + + private static final UUID COMPANY_A = + UUID.fromString("a1000000-0000-0000-0000-000000000001"); + private static final UUID COMPANY_B = + UUID.fromString("b1000000-0000-0000-0000-000000000002"); + private static final UUID ADMIN_A = + UUID.fromString("a2000000-0000-0000-0000-000000000001"); + private static final UUID HR_A = + UUID.fromString("a2000000-0000-0000-0000-000000000002"); + private static final UUID ADMIN_B = + UUID.fromString("b2000000-0000-0000-0000-000000000001"); + private static final UUID EVENT_A = + UUID.fromString("a3000000-0000-0000-0000-000000000001"); + private static final String ADMIN_A_EMAIL = "outbox.admin.a@example.com"; + private static final String HR_A_EMAIL = "outbox.hr.a@example.com"; + private static final String ADMIN_B_EMAIL = "outbox.admin.b@example.com"; + private static final String PASSWORD = "Test-password-1!"; + + @LocalServerPort + private int port; + + @Autowired + private JdbcTemplate jdbcTemplate; + + @Autowired + private PasswordEncoder passwordEncoder; + + @Autowired + private OutboxProcessor outboxProcessor; + + @Autowired + private ManualRetryProbeHandler probeHandler; + + private final HttpClient httpClient = HttpClient.newHttpClient(); + + @BeforeEach + void resetAndSeed() { + probeHandler.reset(); + jdbcTemplate.update("DELETE FROM outbox_manual_retry"); + jdbcTemplate.update("DELETE FROM ai_candidate_decision_task"); + jdbcTemplate.update("DELETE FROM ai_candidate_decision"); + jdbcTemplate.update("DELETE FROM ai_candidate_decision_batch"); + jdbcTemplate.update("DELETE FROM ai_candidate"); + jdbcTemplate.update("DELETE FROM ai_question"); + jdbcTemplate.update("DELETE FROM ai_attempt"); + jdbcTemplate.update("DELETE FROM ai_run"); + jdbcTemplate.update("DELETE FROM event_consumption"); + jdbcTemplate.update("DELETE FROM event_publication"); + jdbcTemplate.update("DELETE FROM audit_event"); + jdbcTemplate.update("DELETE FROM worker_response_upload"); + jdbcTemplate.update("DELETE FROM worker_document_upload_idempotency"); + jdbcTemplate.update("DELETE FROM worker_response"); + jdbcTemplate.update("DELETE FROM worker_link"); + jdbcTemplate.update("DELETE FROM document_request_draft_type"); + jdbcTemplate.update("DELETE FROM document_request_draft"); + jdbcTemplate.update("DELETE FROM task_evidence"); + jdbcTemplate.update("DELETE FROM external_submission"); + jdbcTemplate.update("DELETE FROM approval_request"); + jdbcTemplate.update("DELETE FROM task_transition_history"); + jdbcTemplate.update("DELETE FROM task_checklist_item"); + jdbcTemplate.update("DELETE FROM stored_file"); + jdbcTemplate.update("DELETE FROM worker_document"); + jdbcTemplate.update("DELETE FROM task"); + jdbcTemplate.update("DELETE FROM workflow_case"); + jdbcTemplate.update("DELETE FROM worker"); + jdbcTemplate.update("DELETE FROM refresh_token"); + jdbcTemplate.update("DELETE FROM user_account"); + jdbcTemplate.update("DELETE FROM company"); + + insertCompany(COMPANY_A, "Outbox 사업장 A"); + insertCompany(COMPANY_B, "Outbox 사업장 B"); + String passwordHash = passwordEncoder.encode(PASSWORD); + insertUser(ADMIN_A, COMPANY_A, ADMIN_A_EMAIL, "ADMIN", passwordHash); + insertUser(HR_A, COMPANY_A, HR_A_EMAIL, "HR", passwordHash); + insertUser(ADMIN_B, COMPANY_B, ADMIN_B_EMAIL, "ADMIN", passwordHash); + insertReviewRequiredEvent(EVENT_A, COMPANY_A); + } + + @Test + void exhaustedEventGetsOneFreshHandlerAttemptAfterManualRetry() throws Exception { + String token = login(ADMIN_A_EMAIL); + + HttpResponse response = retry( + EVENT_A, + token, + "outbox-retry-exhausted", + validBody(0) + ); + + assertThat(response.statusCode()).isEqualTo(202); + assertThat(attemptCount(EVENT_A)).isZero(); + assertThat(jdbcTemplate.queryForObject( + "SELECT previous_attempt_count FROM outbox_manual_retry WHERE event_id = ?", + Integer.class, + EVENT_A + )).isEqualTo(8); + + assertThat(outboxProcessor.processAvailable()).isEqualTo(1); + assertThat(probeHandler.invocationCount()).isEqualTo(1); + assertThat(status(EVENT_A)).isEqualTo("COMPLETED"); + assertThat(attemptCount(EVENT_A)).isEqualTo(1); + } + + @Test + void adminRetriesOnceWithoutExposingPayloadAndDuplicateRequestIsReplayed() throws Exception { + String token = login(ADMIN_A_EMAIL); + String body = """ + {"expected_version":0,"reason":"내부 handler 복구와 점검을 완료했습니다."} + """; + + HttpResponse first = retry(EVENT_A, token, "outbox-retry-001", body); + + assertThat(first.statusCode()).isEqualTo(202); + assertThat(JsonPath.read(first.body(), "$.accepted_status")).isEqualTo("PENDING"); + assertThat(JsonPath.read(first.body(), "$.already_requested")).isFalse(); + assertThat(first.body()) + .doesNotContain("payload_json", "last_error_code", "secret-value"); + assertThat(status(EVENT_A)).isEqualTo("PENDING"); + assertThat(lastErrorCode(EVENT_A)).isNull(); + assertThat(count("outbox_manual_retry")).isEqualTo(1); + assertThat(count("audit_event")).isEqualTo(1); + assertThat(jdbcTemplate.queryForObject( + "SELECT action FROM audit_event WHERE target_id = ?", + String.class, + EVENT_A + )).isEqualTo("OUTBOX_MANUAL_RETRY_REQUESTED"); + + HttpResponse duplicate = retry(EVENT_A, token, "outbox-retry-001", body); + + assertThat(duplicate.statusCode()).isEqualTo(202); + assertThat(JsonPath.read(duplicate.body(), "$.already_requested")).isTrue(); + assertThat(count("outbox_manual_retry")).isEqualTo(1); + assertThat(count("audit_event")).isEqualTo(1); + + HttpResponse conflictingReuse = retry( + EVENT_A, + token, + "outbox-retry-001", + """ + {"expected_version":0,"reason":"다른 원인으로 재처리를 다시 요청합니다."} + """ + ); + assertThat(conflictingReuse.statusCode()).isEqualTo(409); + } + + @Test + void roleTenantStateAndVersionAreChecked() throws Exception { + String adminA = login(ADMIN_A_EMAIL); + + assertThat(retry( + EVENT_A, + login(HR_A_EMAIL), + "outbox-retry-hr", + validBody(0) + ).statusCode()).isEqualTo(403); + assertThat(retry( + EVENT_A, + login(ADMIN_B_EMAIL), + "outbox-retry-other-company", + validBody(0) + ).statusCode()).isEqualTo(404); + assertThat(retry( + EVENT_A, + adminA, + "outbox-retry-stale", + validBody(1) + ).statusCode()).isEqualTo(409); + + jdbcTemplate.update( + "UPDATE event_publication SET status = 'COMPLETED', " + + "last_error_code = NULL, completed_at = CURRENT_TIMESTAMP WHERE event_id = ?", + EVENT_A + ); + assertThat(retry( + EVENT_A, + adminA, + "outbox-retry-completed", + validBody(0) + ).statusCode()).isEqualTo(409); + } + + @Test + void expectedVersionIsRequired() throws Exception { + String token = login(ADMIN_A_EMAIL); + + assertThat(retry( + EVENT_A, + token, + "outbox-retry-missing-version", + """ + {"reason":"내부 handler 복구와 점검을 완료했습니다."} + """ + ).statusCode()).isEqualTo(400); + assertThat(retry( + EVENT_A, + token, + "outbox-retry-null-version", + """ + {"expected_version":null,"reason":"내부 handler 복구와 점검을 완료했습니다."} + """ + ).statusCode()).isEqualTo(400); + assertThat(status(EVENT_A)).isEqualTo("REVIEW_REQUIRED"); + assertThat(count("outbox_manual_retry")).isZero(); + } + + @Test + void idempotencyKeyHeaderIsRequiredWithoutChangingEvent() throws Exception { + String token = login(ADMIN_A_EMAIL); + + HttpResponse response = retryWithoutIdempotencyKey( + EVENT_A, + token, + validBody(0) + ); + + assertThat(response.statusCode()).isEqualTo(400); + assertThat(JsonPath.read(response.body(), "$.code")).isEqualTo("INVALID_REQUEST"); + assertThat(status(EVENT_A)).isEqualTo("REVIEW_REQUIRED"); + assertThat(attemptCount(EVENT_A)).isEqualTo(8); + assertThat(count("outbox_manual_retry")).isZero(); + assertThat(count("audit_event")).isZero(); + } + + @Test + void concurrentRequestsAllowOnlyOneRetry() throws Exception { + String token = login(ADMIN_A_EMAIL); + ExecutorService executor = Executors.newFixedThreadPool(2); + CountDownLatch ready = new CountDownLatch(2); + CountDownLatch start = new CountDownLatch(1); + try { + Future> first = executor.submit( + () -> retryAfterSignal("outbox-concurrent-1", token, ready, start) + ); + Future> second = executor.submit( + () -> retryAfterSignal("outbox-concurrent-2", token, ready, start) + ); + assertThat(ready.await(5, TimeUnit.SECONDS)).isTrue(); + start.countDown(); + + assertThat(List.of( + first.get(10, TimeUnit.SECONDS).statusCode(), + second.get(10, TimeUnit.SECONDS).statusCode() + )).containsExactlyInAnyOrder(202, 409); + } finally { + start.countDown(); + executor.shutdownNow(); + } + + assertThat(count("outbox_manual_retry")).isEqualTo(1); + assertThat(count("audit_event")).isEqualTo(1); + assertThat(status(EVENT_A)).isEqualTo("PENDING"); + } + + @Test + void openApiPublishesAdminRetryContract() throws Exception { + HttpResponse response = httpClient.send( + HttpRequest.newBuilder(uri("/v3/api-docs")).GET().build(), + HttpResponse.BodyHandlers.ofString() + ); + + assertThat(response.statusCode()).isEqualTo(200); + assertThat(JsonPath.read( + response.body(), + "$.paths['/api/v1/admin/outbox-events/{eventId}/retry'].post.operationId" + )).isEqualTo("retryOutboxEvent"); + assertThat(response.body()) + .contains("Idempotency-Key", "expected_version", "reason", "bearerAuth"); + } + + private HttpResponse retryAfterSignal( + String key, + String token, + CountDownLatch ready, + CountDownLatch start + ) throws Exception { + ready.countDown(); + start.await(5, TimeUnit.SECONDS); + return retry(EVENT_A, token, key, validBody(0)); + } + + private HttpResponse retry(UUID eventId, String token, String key, String body) + throws Exception { + return httpClient.send( + HttpRequest.newBuilder(uri("/api/v1/admin/outbox-events/" + eventId + "/retry")) + .header(HttpHeaders.AUTHORIZATION, "Bearer " + token) + .header(HttpHeaders.CONTENT_TYPE, "application/json") + .header("Idempotency-Key", key) + .POST(HttpRequest.BodyPublishers.ofString(body)) + .build(), + HttpResponse.BodyHandlers.ofString() + ); + } + + private HttpResponse retryWithoutIdempotencyKey(UUID eventId, String token, String body) + throws Exception { + return httpClient.send( + HttpRequest.newBuilder(uri("/api/v1/admin/outbox-events/" + eventId + "/retry")) + .header(HttpHeaders.AUTHORIZATION, "Bearer " + token) + .header(HttpHeaders.CONTENT_TYPE, "application/json") + .POST(HttpRequest.BodyPublishers.ofString(body)) + .build(), + HttpResponse.BodyHandlers.ofString() + ); + } + + private String validBody(long version) { + return """ + {"expected_version":%d,"reason":"내부 handler 복구와 점검을 완료했습니다."} + """.formatted(version); + } + + private String login(String email) throws Exception { + HttpResponse response = httpClient.send( + HttpRequest.newBuilder(uri("/api/v1/auth/login")) + .header(HttpHeaders.CONTENT_TYPE, "application/json") + .POST(HttpRequest.BodyPublishers.ofString(""" + {"email":"%s","password":"%s"} + """.formatted(email, PASSWORD))) + .build(), + HttpResponse.BodyHandlers.ofString() + ); + assertThat(response.statusCode()).isEqualTo(200); + return JsonPath.read(response.body(), "$.access_token"); + } + + private URI uri(String path) { + return URI.create("http://localhost:" + port + path); + } + + private void insertCompany(UUID companyId, String name) { + jdbcTemplate.update( + "INSERT INTO company (company_id, name, status) VALUES (?, ?, 'ACTIVE')", + companyId, + name + ); + } + + private void insertUser( + UUID userId, + UUID companyId, + String email, + String role, + String passwordHash + ) { + jdbcTemplate.update( + """ + INSERT INTO user_account ( + user_id, company_id, email, normalized_email, + password_hash, role, status, display_name + ) VALUES (?, ?, ?, ?, ?, ?, 'ACTIVE', '운영 테스트') + """, + userId, + companyId, + email, + email, + passwordHash, + role + ); + } + + private void insertReviewRequiredEvent(UUID eventId, UUID companyId) { + jdbcTemplate.update( + """ + INSERT INTO event_publication ( + event_id, company_id, event_type, payload_version, + aggregate_type, aggregate_id, actor_type, actor_id, + request_id, payload_json, status, attempt_count, + last_error_code, occurred_at, created_at, updated_at, version + ) VALUES ( + ?, ?, 'ReliabilityTestRequested', '1', + 'ReliabilityProbe', ?, 'SYSTEM_RULE', NULL, + 'outbox-manual-retry-fixture', '{"result":"secret-value"}', + 'REVIEW_REQUIRED', 8, 'TEST_PAYLOAD_REJECTED', + CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, 0 + ) + """, + eventId, + companyId, + UUID.randomUUID() + ); + } + + private int count(String table) { + return jdbcTemplate.queryForObject("SELECT COUNT(*) FROM " + table, Integer.class); + } + + private String status(UUID eventId) { + return jdbcTemplate.queryForObject( + "SELECT status FROM event_publication WHERE event_id = ?", + String.class, + eventId + ); + } + + private String lastErrorCode(UUID eventId) { + return jdbcTemplate.queryForObject( + "SELECT last_error_code FROM event_publication WHERE event_id = ?", + String.class, + eventId + ); + } + + private int attemptCount(UUID eventId) { + return jdbcTemplate.queryForObject( + "SELECT attempt_count FROM event_publication WHERE event_id = ?", + Integer.class, + eventId + ); + } + + @TestConfiguration(proxyBeanMethods = false) + static class ManualRetryTestConfiguration { + + @Bean + ManualRetryProbeHandler manualRetryProbeHandler() { + return new ManualRetryProbeHandler(); + } + } + + static final class ManualRetryProbeHandler implements DomainEventHandler { + + private final AtomicInteger invocationCount = new AtomicInteger(); + + @Override + public String handlerName() { + return "manual-retry-probe-v1"; + } + + @Override + public boolean supports(String eventType) { + return "ReliabilityTestRequested".equals(eventType); + } + + @Override + public void handle(DomainEventEnvelope event) { + invocationCount.incrementAndGet(); + } + + int invocationCount() { + return invocationCount.get(); + } + + void reset() { + invocationCount.set(0); + } + } +}