diff --git a/README.md b/README.md index aa50c290..9cbe510e 100644 --- a/README.md +++ b/README.md @@ -277,6 +277,7 @@ src/main/java/com/fowoco/server/ | Agent DB 정보 보충 | [Slot 조회·재호출](docs/ai-slot-resolution.md) | canonical key allow-list, tenant 조회와 ANALYZE 재호출 기준 | | AI 단계별 성능 측정 | [AI 파이프라인 관측·Prometheus 가이드](docs/ai-pipeline-observability.md) | PLAN·Slot·ANALYZE·Renewal 구간의 정량 평가와 로컬 Prometheus 확인 기준 | | 이벤트 유실·재처리 | [Outbox 운영 가이드](docs/reliability/transactional-outbox.md) | 이벤트 발행, lease, 재시도와 장애 복구 기준 | +| 파일 rollback·orphan 대응 | [File Storage rollback 보상 운영 가이드](docs/reliability/file-storage-rollback-compensation.md) | atomic finalize, rollback cleanup, `UNKNOWN` reconciliation과 배포 volume Smoke 기준 | | 구현 계획·업무 상태 | [Server Roadmap](https://github.com/orgs/fowoco/projects/3) · [Issues](https://github.com/fowoco/server/issues) | 실제 담당자, 우선순위와 진행 상태 | | 전체 설명·운영 가이드 | [Server Wiki](https://github.com/fowoco/server/wiki) | 초보자용 아키텍처·API·배포 설명 | diff --git a/docs/deployment-runbook.md b/docs/deployment-runbook.md index 891cf55e..460781c5 100644 --- a/docs/deployment-runbook.md +++ b/docs/deployment-runbook.md @@ -158,10 +158,17 @@ Seed의 수량과 고정 ID도 첫 기동과 같아야 합니다. 4. 로그인과 타 사업장 접근 차단 확인 5. `POST /api/v1/ai-runs`의 실제 Server→AI 왕복 확인 6. 후보 채택 후 Case·Task 조회 확인 -7. Worker Link 대표 흐름 확인 +7. Worker Link 대표 흐름과 동일 `Idempotency-Key` 문서 재시도가 같은 `upload_id`로 + 수렴하는지 확인 8. SMS가 활성화된 환경에서는 실제 수신·링크 접속·중복 발송 방지 확인 9. SMTP가 활성화된 환경에서는 재설정 메일 수신·링크 token·새 비밀번호 로그인 확인 10. 로그에서 AiRun·Renewal `TOTAL` 단계와 안전한 `failure_code`가 기록되는지 확인 +11. 실제 `FILE_STORAGE_LOCAL_PATH` volume에 최종 파일만 한 개 있고 + `.fowoco-upload-*.tmp`가 남지 않았는지 확인 + +파일 rollback cleanup 실패나 transaction `UNKNOWN` 로그가 있으면 +[File Storage rollback 보상 운영 가이드](reliability/file-storage-rollback-compensation.md)의 +DB·volume 대조 절차를 따른다. Runtime 장애 테스트에서는 가짜 AI 결과를 만들지 않고 안전한 오류 또는 수동 처리 상태로 남아야 합니다. @@ -180,6 +187,15 @@ kubectl -n fowoco get events --sort-by=.lastTimestamp 애플리케이션 문제이며 DB migration이 이전 이미지와 호환될 때만 승인 후 이전 image SHA로 되돌립니다. +Worker Link 문서 멱등성 V51을 적용한 뒤 이전 Server image로 되돌릴 때는 schema 호환과 +멱등성 의미 호환을 구분합니다. 이전 image는 nullable hash column을 무시하고 기동할 수 +있지만, 새 version이 `canonical:`로 기록한 성공 결과를 기존 +`clientRequestId` 조회로 재사용하지 못합니다. rollback 가능한 기간에는 Client가 +`Idempotency-Key`와 multipart `clientRequestId`를 함께 보내는 동작을 유지하고, rollback +후에는 새 version에서 성공한 문서 업로드를 자동 재시도하지 않습니다. 상세 점검은 +[File Storage rollback 보상 운영 가이드](reliability/file-storage-rollback-compensation.md)의 +Rollback 원칙을 따릅니다. + ```bash kubectl -n fowoco set image deployment/server server=ghcr.io/fowoco/server: kubectl -n fowoco rollout status deployment/server --timeout=180s diff --git a/docs/reliability/file-storage-rollback-compensation.md b/docs/reliability/file-storage-rollback-compensation.md new file mode 100644 index 00000000..933810fd --- /dev/null +++ b/docs/reliability/file-storage-rollback-compensation.md @@ -0,0 +1,105 @@ +# File Storage rollback 보상 운영 가이드 + +## 목적과 보장 범위 + +로컬 파일시스템은 PostgreSQL 트랜잭션에 참여하지 않는다. 파일을 먼저 저장한 뒤 +`stored_file` 또는 감사로그 영속화가 실패하면 DB는 rollback되지만 파일만 남을 수 있다. +Server는 파일 저장 전에 transaction synchronization을 등록하고, transaction 결과가 +`ROLLED_BACK`일 때 서버가 생성한 `storage_key`를 `deleteIfExists`로 정리한다. + +현재 보장 범위는 다음과 같다. + +- 최종 파일은 storage root 안의 임시 파일에 완전히 기록한 후 atomic move로 공개한다. +- 배포 volume이 atomic move를 지원하지 않으면 일반 move로 대체하지 않고 저장을 + 실패시킨다. 불완전한 최종 파일을 노출하지 않는 것이 저장 성공보다 우선한다. +- DB rollback callback은 같은 `storage_key`를 여러 번 정리해도 성공하도록 멱등 삭제한다. +- `FileService` 업로드와 Worker Link 문서 업로드 모두 같은 rollback 보상을 사용한다. +- Worker Link 재시도 key는 원문 대신 SHA-256 hash만 DB에 저장한다. 같은 key와 같은 + 요청은 기존 `upload_id`로 수렴하고, 같은 key의 다른 요청은 + `IDEMPOTENCY_CONFLICT`로 거부한다. + +이 보상은 DB와 파일시스템을 하나의 원자적 transaction으로 바꾸지 않는다. 프로세스가 +파일 finalize 직후 강제 종료되거나 transaction 완료 결과가 `UNKNOWN`이면 운영 확인이 +필요하다. + +## 구조화 로그 + +rollback 보상은 다음 event를 남긴다. + +| event | status | 의미 | 운영 조치 | +| --- | --- | --- | --- | +| `file_storage_cleanup` | `ATTEMPT` | rollback 파일 삭제 시작 | 뒤따르는 동일 `request_id`·`storage_key` 결과 확인 | +| `file_storage_cleanup` | `SUCCEEDED` | 파일이 삭제됐거나 이미 없음 | 별도 조치 없음 | +| `file_storage_cleanup` | `FAILED` | 삭제 중 예외 발생 | DB와 파일을 대조해 orphan 여부 확인 | +| `file_storage_transaction_completion` | `UNKNOWN` | commit/rollback 결과를 확정하지 못함 | 자동 삭제 금지, reconciliation 수행 | + +로그에는 `request_id`, `action`, `storage`, `phase`, 서버 생성 `storage_key`만 사용한다. +원본 Worker Link token, `Idempotency-Key`, 파일명, 파일 내용과 사용자 개인정보를 +추가하지 않는다. + +우선 확인할 검색 조건은 다음과 같다. + +```text +event=file_storage_cleanup status=FAILED +event=file_storage_transaction_completion status=UNKNOWN reconciliation_required=true +``` + +## FAILED와 UNKNOWN 대응 + +`FAILED` 또는 `UNKNOWN` 한 건마다 로그의 `request_id`와 `storage_key`를 기준으로 +`stored_file.storage_key` 행과 실제 storage volume의 최종 파일을 대조한다. + +| DB 행 | 최종 파일 | 판정과 조치 | +| --- | --- | --- | +| 있음 | 있음 | commit된 정상 파일이다. 삭제하지 않는다. | +| 없음 | 있음 | orphan 후보이다. 동일 key가 서버 생성 UUID인지, 연결된 업무 행이 없는지 재확인한 뒤 승인된 운영 절차로 파일만 멱등 삭제한다. | +| 있음 | 없음 | DB가 가리키는 파일이 유실된 상태다. 자동 DB 삭제를 금지하고 복구 또는 재업로드를 결정한다. | +| 없음 | 없음 | rollback 정리가 완료된 상태다. | + +`UNKNOWN`은 실제로 commit됐을 가능성이 있으므로 파일부터 지우지 않는다. 확인 중에는 +원본 파일명이나 token을 로그에 복사하지 않고, 접근이 제한된 DB와 volume에서 서버 생성 +UUID key만 사용한다. 수동 삭제를 재시도한 뒤에는 같은 key의 파일 부재와 DB 행 부재를 +다시 확인하고 incident 기록에 `request_id`와 판정만 남긴다. + +## 배포 전후 Smoke + +`FILE_STORAGE_LOCAL_PATH`가 실제 배포 Pod에 mount된 경로인지 먼저 확인한다. 임시 +container filesystem을 가리키거나 여러 Pod가 서로 다른 local volume을 사용하면 이 +구현의 범위를 벗어난다. + +1. 실제 mount에서 허용 MIME 파일을 업로드하고 다운로드 내용이 같은지 확인한다. +2. Worker Link 문서 제출에 동일한 `Idempotency-Key`와 동일 payload를 두 번 보내 두 + 응답의 `upload_id`가 같은지 확인한다. +3. 같은 key로 내용만 바꾼 요청이 `409 IDEMPOTENCY_CONFLICT`인지 확인한다. +4. DB에는 key별 idempotency 행과 `stored_file` 행이 각각 한 건이고, volume에는 대응하는 + 최종 파일이 한 개뿐인지 확인한다. +5. storage root에 `.fowoco-upload-*.tmp`가 남지 않았는지 확인한다. +6. 검증 환경에서 파일 저장 이후 transaction rollback을 강제로 발생시켜 DB 행, 감사로그, + 최종 파일과 임시 파일이 모두 남지 않는지 확인한다. +7. mount가 atomic move를 지원하지 않으면 배포를 중단한다. 일반 move fallback을 + 추가하지 말고 atomic rename이 가능한 volume 또는 별도 object storage 구현을 선택한다. + +PostgreSQL 16 동시성 검증은 +`WorkerLinkDocumentPostgreSqlIntegrationTest`가 두 HTTP 요청을 Worker Link 행 잠금 +직전에 겹치게 한 뒤 DB 행과 실제 LocalFileStorage 파일이 하나로 수렴하는지 반복 +확인한다. 환경변수가 없으면 테스트가 skip되므로 결과에서 실제 실행 건수를 반드시 +확인한다. + +## Rollback 원칙 + +- 적용된 Flyway migration을 수정하거나 schema history를 조작하지 않는다. +- V51의 hash column은 nullable이므로 이전 image가 같은 schema에서 기동하고 legacy 행을 + 기록할 수 있다. 이는 schema 하위 호환을 의미하며, 신·구 version 사이의 멱등성 결과 + 재사용까지 보장한다는 뜻은 아니다. +- 새 version은 `client_request_id`에 `canonical:`를 기록하지만 이전 + version은 multipart `clientRequestId` 원문으로 기존 결과를 조회한다. 따라서 새 version이 + 성공시킨 업로드를 이전 image로 rollback한 뒤 재시도하면 기존 결과를 찾지 못하고 중복 + 업로드할 수 있다. +- rollback 가능한 배포 기간에는 Client가 `Idempotency-Key`와 `clientRequestId`를 함께 + 보내는 현재 동작을 유지한다. 이전 Server version을 지원하지 않기로 확정하기 전에는 + `clientRequestId` 전송을 제거하지 않는다. +- 이전 image로 rollback한 경우 새 version에서 성공한 Worker Link 문서 요청의 자동 재시도를 + 피하고, 재시도가 발생했다면 `worker_document_upload_idempotency`, `stored_file`과 실제 + volume을 대조해 중복 여부를 확인한다. +- 이전 image로 되돌린 뒤에도 orphan 후보는 위 reconciliation 표로 판단한다. +- cleanup 실패를 숨기기 위해 `stored_file` 행이나 파일을 일괄 삭제하지 않는다. diff --git a/src/main/java/com/fowoco/server/file/application/FileService.java b/src/main/java/com/fowoco/server/file/application/FileService.java index a8f13d86..0451b472 100644 --- a/src/main/java/com/fowoco/server/file/application/FileService.java +++ b/src/main/java/com/fowoco/server/file/application/FileService.java @@ -47,11 +47,13 @@ public class FileService { ); private static final String HWP_EXTENSION = ".hwp"; private static final String HWPX_EXTENSION = ".hwpx"; + private static final String FILE_UPLOAD_ACTION = "file_upload"; private final StoredFileRepository storedFileRepository; private final HwpSignatureValidator hwpSignatureValidator; private final HwpxSignatureValidator hwpxSignatureValidator; private final FileStorage fileStorage; + private final FileStorageRollbackCompensation rollbackCompensation; private final TaskRepository taskRepository; private final WorkerRepository workerRepository; private final AuditEventRepository auditRepository; @@ -64,6 +66,7 @@ public FileService( HwpSignatureValidator hwpSignatureValidator, HwpxSignatureValidator hwpxSignatureValidator, FileStorage fileStorage, + FileStorageRollbackCompensation rollbackCompensation, TaskRepository taskRepository, WorkerRepository workerRepository, AuditEventRepository auditRepository, @@ -75,6 +78,7 @@ public FileService( this.hwpSignatureValidator = hwpSignatureValidator; this.hwpxSignatureValidator = hwpxSignatureValidator; this.fileStorage = fileStorage; + this.rollbackCompensation = rollbackCompensation; this.taskRepository = taskRepository; this.workerRepository = workerRepository; this.auditRepository = auditRepository; @@ -128,7 +132,10 @@ public StoredFile upload(FileCreateCommand command, ActorContext actor, RequestM now ); + FileStorageRollbackCompensation.Registration rollbackRegistration = + rollbackCompensation.register(storageKey, metadata, FILE_UPLOAD_ACTION); fileStorage.store(storageKey, new java.io.ByteArrayInputStream(contentBytes), command.size(), command.mimeType()); + rollbackRegistration.markCreated(); storedFileRepository.insert(storedFile); appendAudit( diff --git a/src/main/java/com/fowoco/server/file/application/FileStorageRollbackCompensation.java b/src/main/java/com/fowoco/server/file/application/FileStorageRollbackCompensation.java new file mode 100644 index 00000000..60db33d9 --- /dev/null +++ b/src/main/java/com/fowoco/server/file/application/FileStorageRollbackCompensation.java @@ -0,0 +1,132 @@ +package com.fowoco.server.file.application; + +import com.fowoco.server.common.web.RequestMetadata; +import com.fowoco.server.file.application.port.FileStorage; +import java.util.Objects; +import java.util.concurrent.atomic.AtomicBoolean; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Component; +import org.springframework.transaction.support.TransactionSynchronization; +import org.springframework.transaction.support.TransactionSynchronizationManager; + +@Component +public class FileStorageRollbackCompensation { + + private static final Logger log = LoggerFactory.getLogger(FileStorageRollbackCompensation.class); + private static final String COMPLETION_PHASE = "transaction_after_completion"; + + private final FileStorage fileStorage; + + public FileStorageRollbackCompensation(FileStorage fileStorage) { + this.fileStorage = fileStorage; + } + + public Registration register(String storageKey, RequestMetadata metadata, String action) { + requireText(storageKey, "storageKey"); + Objects.requireNonNull(metadata, "metadata must not be null"); + requireText(action, "action"); + if (!TransactionSynchronizationManager.isActualTransactionActive() + || !TransactionSynchronizationManager.isSynchronizationActive()) { + throw new IllegalStateException( + "File storage rollback compensation requires an active transaction synchronization." + ); + } + + Registration registration = new Registration(); + TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { + @Override + public void afterCompletion(int status) { + handleCompletion(status, storageKey, metadata.requestId(), action, registration.fileCreated()); + } + }); + return registration; + } + + private void handleCompletion( + int status, + String storageKey, + String requestId, + String action, + boolean fileCreated + ) { + if (!fileCreated) { + return; + } + if (status == TransactionSynchronization.STATUS_COMMITTED) { + return; + } + if (status == TransactionSynchronization.STATUS_ROLLED_BACK) { + cleanup(storageKey, requestId, action); + return; + } + log.warn( + "event=file_storage_transaction_completion status=UNKNOWN request_id={} action={} " + + "storage={} phase={} storage_key={} reconciliation_required=true completion_status={}", + requestId, + action, + storageName(), + COMPLETION_PHASE, + storageKey, + status + ); + } + + private void cleanup(String storageKey, String requestId, String action) { + log.info( + "event=file_storage_cleanup status=ATTEMPT request_id={} action={} " + + "storage={} phase={} storage_key={}", + requestId, + action, + storageName(), + COMPLETION_PHASE, + storageKey + ); + try { + fileStorage.deleteIfExists(storageKey); + log.info( + "event=file_storage_cleanup status=SUCCEEDED request_id={} action={} " + + "storage={} phase={} storage_key={}", + requestId, + action, + storageName(), + COMPLETION_PHASE, + storageKey + ); + } catch (RuntimeException exception) { + log.warn( + "event=file_storage_cleanup status=FAILED request_id={} action={} " + + "storage={} phase={} storage_key={}", + requestId, + action, + storageName(), + COMPLETION_PHASE, + storageKey, + exception + ); + } + } + + private String storageName() { + return fileStorage.getClass().getSimpleName(); + } + + private void requireText(String value, String fieldName) { + if (value == null || value.isBlank()) { + throw new IllegalArgumentException(fieldName + " must not be blank"); + } + } + + public static final class Registration { + + private final AtomicBoolean fileCreated = new AtomicBoolean(); + + public void markCreated() { + fileCreated.set(true); + } + + private boolean fileCreated() { + return fileCreated.get(); + } + } +} diff --git a/src/main/java/com/fowoco/server/file/application/port/FileStorage.java b/src/main/java/com/fowoco/server/file/application/port/FileStorage.java index 63298ec5..f2d35afd 100644 --- a/src/main/java/com/fowoco/server/file/application/port/FileStorage.java +++ b/src/main/java/com/fowoco/server/file/application/port/FileStorage.java @@ -8,4 +8,6 @@ public interface FileStorage { void store(String storageKey, InputStream content, long size, String mimeType); Optional open(String storageKey); + + void deleteIfExists(String storageKey); } diff --git a/src/main/java/com/fowoco/server/file/infrastructure/LocalFileStorage.java b/src/main/java/com/fowoco/server/file/infrastructure/LocalFileStorage.java index 8ad3a8ca..9ba9e436 100644 --- a/src/main/java/com/fowoco/server/file/infrastructure/LocalFileStorage.java +++ b/src/main/java/com/fowoco/server/file/infrastructure/LocalFileStorage.java @@ -4,8 +4,11 @@ import java.io.IOException; import java.io.InputStream; import java.io.UncheckedIOException; +import java.nio.file.FileAlreadyExistsException; import java.nio.file.Files; +import java.nio.file.LinkOption; import java.nio.file.Path; +import java.nio.file.StandardCopyOption; import java.util.Optional; import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Component; @@ -13,6 +16,9 @@ @Component public class LocalFileStorage implements FileStorage { + private static final String TEMPORARY_FILE_PREFIX = ".fowoco-upload-"; + private static final String TEMPORARY_FILE_SUFFIX = ".tmp"; + private final Path rootDirectory; public LocalFileStorage(@Value("${app.file-storage.local-path}") String localPath) { @@ -21,12 +27,34 @@ public LocalFileStorage(@Value("${app.file-storage.local-path}") String localPat @Override public void store(String storageKey, InputStream content, long size, String mimeType) { + Path temporary = null; try { Files.createDirectories(rootDirectory); Path target = resolveSafely(storageKey); - Files.copy(content, target); + if (Files.exists(target, LinkOption.NOFOLLOW_LINKS)) { + throw new FileAlreadyExistsException(target.toString()); + } + + temporary = Files.createTempFile( + rootDirectory, + TEMPORARY_FILE_PREFIX, + TEMPORARY_FILE_SUFFIX + ); + Files.copy(content, temporary, StandardCopyOption.REPLACE_EXISTING); + long storedSize = Files.size(temporary); + if (storedSize != size) { + throw new IOException( + "stored file size does not match the declared size: expected=" + + size + ", actual=" + storedSize + ); + } + moveToFinalPath(temporary, target); } catch (IOException exception) { + deleteTemporaryAfterFailure(temporary, exception); throw new UncheckedIOException("failed to store file: " + storageKey, exception); + } catch (RuntimeException exception) { + deleteTemporaryAfterFailure(temporary, exception); + throw exception; } } @@ -43,6 +71,31 @@ public Optional open(String storageKey) { } } + @Override + public void deleteIfExists(String storageKey) { + Path target = resolveSafely(storageKey); + try { + Files.deleteIfExists(target); + } catch (IOException exception) { + throw new UncheckedIOException("failed to delete file: " + storageKey, exception); + } + } + + private void moveToFinalPath(Path temporary, Path target) throws IOException { + Files.move(temporary, target, StandardCopyOption.ATOMIC_MOVE); + } + + private void deleteTemporaryAfterFailure(Path temporary, Throwable failure) { + if (temporary == null) { + return; + } + try { + Files.deleteIfExists(temporary); + } catch (IOException cleanupFailure) { + failure.addSuppressed(cleanupFailure); + } + } + private Path resolveSafely(String storageKey) { Path target = rootDirectory.resolve(storageKey).normalize(); if (!target.startsWith(rootDirectory)) { diff --git a/src/main/java/com/fowoco/server/workerlink/api/WorkerLinkDocumentController.java b/src/main/java/com/fowoco/server/workerlink/api/WorkerLinkDocumentController.java index dde6990a..24911e25 100644 --- a/src/main/java/com/fowoco/server/workerlink/api/WorkerLinkDocumentController.java +++ b/src/main/java/com/fowoco/server/workerlink/api/WorkerLinkDocumentController.java @@ -14,6 +14,9 @@ import io.swagger.v3.oas.annotations.responses.ApiResponses; import io.swagger.v3.oas.annotations.tags.Tag; import jakarta.servlet.http.HttpServletRequest; +import jakarta.validation.constraints.NotBlank; +import jakarta.validation.constraints.Pattern; +import jakarta.validation.constraints.Size; import java.io.IOException; import java.io.UncheckedIOException; import org.springframework.http.HttpStatus; @@ -51,6 +54,7 @@ public WorkerLinkDocumentController(WorkerLinkDocumentService workerLinkDocument ) ), @ApiResponse(responseCode = "400", ref = "#/components/responses/BadRequest"), + @ApiResponse(responseCode = "409", description = "같은 멱등성 키의 요청 내용이 기존 업로드와 다름"), @ApiResponse(responseCode = "410", description = "링크를 찾을 수 없거나 더 이상 사용할 수 없음"), @ApiResponse(responseCode = "413", description = "파일 크기 초과"), @ApiResponse(responseCode = "415", ref = "#/components/responses/UnsupportedMediaType"), @@ -65,9 +69,38 @@ public WorkerLinkDocumentController(WorkerLinkDocumentService workerLinkDocument public ResponseEntity upload( @Parameter(description = "근로자 링크 토큰") @PathVariable String token, @Parameter(description = "업로드할 파일") @RequestParam("file") MultipartFile file, - @Parameter(description = "문서 유형") @RequestParam(value = "documentType", required = false) String documentType, - @Parameter(description = "클라이언트 중복 방지 키") @RequestParam("clientRequestId") String clientRequestId, - @RequestHeader(value = "Idempotency-Key", required = false) String idempotencyKey, + @Parameter( + name = "documentType", + description = "문서 유형", + schema = @Schema( + type = "string", + allowableValues = { + "PASSPORT_COPY", + "ARC", + "CONTRACT", + "PERMIT", + "EMPLOYMENT_EXTENSION_APPLICATION", + "INTEGRATED_APPLICATION", + "RESIDENCE_PROOF" + } + ) + ) + @RequestParam(value = "documentType", required = false) String documentType, + @Parameter( + description = "더 이상 멱등성 판단에 사용하지 않는 이전 클라이언트 요청 식별자", + deprecated = true + ) + @RequestParam(value = "clientRequestId", required = false) String clientRequestId, + @Parameter( + description = "문서 업로드 재시도를 식별하는 필수 키", + required = true, + example = "worker-upload-018f6b65" + ) + @RequestHeader(value = "Idempotency-Key", required = false) + @NotBlank + @Size(min = 8, max = 100) + @Pattern(regexp = "^[A-Za-z0-9][A-Za-z0-9._:-]*$") + String idempotencyKey, HttpServletRequest servletRequest ) { if (file.isEmpty() || file.getOriginalFilename() == null || file.getOriginalFilename().isBlank()) { @@ -81,6 +114,7 @@ public ResponseEntity upload( file.getSize(), documentType, clientRequestId, + idempotencyKey, file.getInputStream() ); WorkerLinkDocumentUploadResult result = workerLinkDocumentService.upload( diff --git a/src/main/java/com/fowoco/server/workerlink/application/WorkerLinkDocumentService.java b/src/main/java/com/fowoco/server/workerlink/application/WorkerLinkDocumentService.java index 92a88dac..1ceff400 100644 --- a/src/main/java/com/fowoco/server/workerlink/application/WorkerLinkDocumentService.java +++ b/src/main/java/com/fowoco/server/workerlink/application/WorkerLinkDocumentService.java @@ -6,24 +6,41 @@ import com.fowoco.server.audit.domain.AuditEvent; import com.fowoco.server.audit.domain.AuditTargetType; 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.file.application.FileStorageRollbackCompensation; import com.fowoco.server.file.application.error.FileErrorCode; import com.fowoco.server.file.application.port.FileStorage; import com.fowoco.server.file.application.port.StoredFileRepository; import com.fowoco.server.file.domain.StoredFile; +import com.fowoco.server.worker.domain.DocumentType; import com.fowoco.server.workerlink.application.error.WorkerLinkErrorCode; +import com.fowoco.server.workerlink.application.port.WorkerDocumentUploadIdempotencyRecord; import com.fowoco.server.workerlink.application.port.WorkerDocumentUploadIdempotencyRepository; import com.fowoco.server.workerlink.application.port.WorkerLinkRepository; import com.fowoco.server.workerlink.application.port.WorkerLinkTenantBootstrap; import com.fowoco.server.workerlink.domain.WorkerLink; import com.fowoco.server.workerlink.infrastructure.security.WorkerLinkHasher; +import java.io.FilterInputStream; +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.io.UncheckedIOException; +import java.nio.ByteBuffer; +import java.nio.charset.StandardCharsets; +import java.security.DigestInputStream; +import java.security.MessageDigest; +import java.security.NoSuchAlgorithmException; import java.time.Clock; import java.time.Instant; +import java.util.HexFormat; +import java.util.Locale; import java.util.Optional; import java.util.Set; import java.util.UUID; +import java.util.regex.Pattern; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; @@ -31,9 +48,16 @@ public class WorkerLinkDocumentService { private static final String AUDIT_EVENT_VERSION = "1"; - - //FileService(#13)와 동일한 값 + private static final String DEFAULT_DOCUMENT_PURPOSE = "WORKER_LINK_SUBMISSION"; + private static final String REQUEST_HASH_VERSION = "worker-link-document-upload:v1"; + private static final String ROLLBACK_ACTION = "worker_link_document_upload"; private static final long MAX_FILE_SIZE_BYTES = 20L * 1024 * 1024; + private static final int MAX_FILE_NAME_LENGTH = 255; + private static final int MIN_IDEMPOTENCY_KEY_LENGTH = 8; + private static final int MAX_IDEMPOTENCY_KEY_LENGTH = 100; + private static final Pattern IDEMPOTENCY_KEY_PATTERN = Pattern.compile( + "^[A-Za-z0-9][A-Za-z0-9._:-]*$" + ); private static final Set ALLOWED_MIME_TYPES = Set.of( "image/jpeg", "image/png", @@ -48,6 +72,7 @@ public class WorkerLinkDocumentService { private final StoredFileRepository storedFileRepository; private final WorkerDocumentUploadIdempotencyRepository uploadIdempotencyRepository; private final FileStorage fileStorage; + private final FileStorageRollbackCompensation rollbackCompensation; private final AuditEventRepository auditRepository; private final UuidGenerator uuidGenerator; private final Clock clock; @@ -60,6 +85,7 @@ public WorkerLinkDocumentService( StoredFileRepository storedFileRepository, WorkerDocumentUploadIdempotencyRepository uploadIdempotencyRepository, FileStorage fileStorage, + FileStorageRollbackCompensation rollbackCompensation, AuditEventRepository auditRepository, UuidGenerator uuidGenerator, Clock clock @@ -71,6 +97,7 @@ public WorkerLinkDocumentService( this.storedFileRepository = storedFileRepository; this.uploadIdempotencyRepository = uploadIdempotencyRepository; this.fileStorage = fileStorage; + this.rollbackCompensation = rollbackCompensation; this.auditRepository = auditRepository; this.uuidGenerator = uuidGenerator; this.clock = clock; @@ -78,14 +105,9 @@ public WorkerLinkDocumentService( @Transactional public WorkerLinkDocumentUploadResult upload(WorkerLinkDocumentUploadCommand command, RequestMetadata metadata) { - if (command.size() > MAX_FILE_SIZE_BYTES) { - throw new ApiException(FileErrorCode.FILE_TOO_LARGE); - } - if (!ALLOWED_MIME_TYPES.contains(command.mimeType())) { - throw new ApiException(FileErrorCode.UNSUPPORTED_FILE_TYPE); - } - + NormalizedUploadRequest request = normalize(command); String tokenHash = workerLinkHasher.hash(command.rawToken()); + String idempotencyKeyHash = workerLinkHasher.hash(request.idempotencyKey()); UUID companyId = workerLinkTenantBootstrap .findCompanyIdByWorkerLinkTokenHash(tokenHash) @@ -93,7 +115,10 @@ public WorkerLinkDocumentUploadResult upload(WorkerLinkDocumentUploadCommand com tenantDatabaseContext.setCompanyIdForCurrentTransaction(companyId); - WorkerLink link = workerLinkRepository.findByTokenHash(tokenHash) + WorkerLink locatedLink = workerLinkRepository.findByTokenHash(tokenHash) + .orElseThrow(() -> new ApiException(WorkerLinkErrorCode.WORKER_LINK_NOT_FOUND)); + WorkerLink link = workerLinkRepository + .findByIdAndCompanyIdForUpdate(locatedLink.workerLinkId(), companyId) .orElseThrow(() -> new ApiException(WorkerLinkErrorCode.WORKER_LINK_NOT_FOUND)); Instant now = clock.instant(); @@ -101,39 +126,38 @@ public WorkerLinkDocumentUploadResult upload(WorkerLinkDocumentUploadCommand com throw new ApiException(WorkerLinkErrorCode.WORKER_LINK_NOT_FOUND); } - Optional existingStoredFileId = uploadIdempotencyRepository - .findStoredFileId(link.workerLinkId(), companyId, command.clientRequestId()); - if (existingStoredFileId.isPresent()) { - StoredFile existingFile = storedFileRepository.findByIdAndCompanyId(existingStoredFileId.get(), companyId) - .orElseThrow(() -> new ApiException(WorkerLinkErrorCode.UPLOAD_NOT_AVAILABLE)); - return new WorkerLinkDocumentUploadResult(existingFile, link.expiresAt()); + Optional existing = uploadIdempotencyRepository + .findByKeyHash(link.workerLinkId(), companyId, idempotencyKeyHash); + if (existing.isPresent()) { + return replayExistingUpload(command, request, link, companyId, existing.get()); } UUID storedFileId = uuidGenerator.generate(); String storageKey = storedFileId.toString(); - String purpose = command.documentType() != null ? command.documentType() : "WORKER_LINK_SUBMISSION"; - - StoredFile storedFile = StoredFile.create( + StoredFile verifiedFile = StoredFile.create( storedFileId, companyId, - command.fileName(), - command.mimeType(), - command.size(), - purpose, + request.fileName(), + request.mimeType(), + request.size(), + request.purpose(), link.taskId(), null, storageKey, now - ); + ).verify(); - StoredFile verifiedFile = storedFile.verify(); + FileStorageRollbackCompensation.Registration rollbackRegistration = + rollbackCompensation.register(storageKey, metadata, ROLLBACK_ACTION); + String contentChecksum = storeAndChecksum(storageKey, command.content(), request, rollbackRegistration); + String requestHash = calculateRequestHash(request, contentChecksum); - fileStorage.store(storageKey, command.content(), command.size(), command.mimeType()); storedFileRepository.insert(verifiedFile); uploadIdempotencyRepository.save( link.workerLinkId(), companyId, - command.clientRequestId(), + idempotencyKeyHash, + requestHash, storedFileId ); @@ -149,10 +173,193 @@ public WorkerLinkDocumentUploadResult upload(WorkerLinkDocumentUploadCommand com metadata.requestId(), metadata.traceId(), AUDIT_EVENT_VERSION, - "근로자 링크로 파일 업로드: " + purpose, + "근로자 링크로 파일 업로드: " + request.purpose(), now )); - return new WorkerLinkDocumentUploadResult(verifiedFile, link.expiresAt()); + return new WorkerLinkDocumentUploadResult(verifiedFile, link.expiresAt()); + } + + private WorkerLinkDocumentUploadResult replayExistingUpload( + WorkerLinkDocumentUploadCommand command, + NormalizedUploadRequest request, + WorkerLink link, + UUID companyId, + WorkerDocumentUploadIdempotencyRecord existing + ) { + String contentChecksum = checksumAndDiscard(command.content(), request.size()); + String requestHash = calculateRequestHash(request, contentChecksum); + if (!existing.requestHash().equals(requestHash)) { + throw new ApiException(WorkerLinkErrorCode.IDEMPOTENCY_CONFLICT); + } + + StoredFile existingFile = storedFileRepository.findByIdAndCompanyId(existing.storedFileId(), companyId) + .orElseThrow(() -> new ApiException(WorkerLinkErrorCode.UPLOAD_NOT_AVAILABLE)); + return new WorkerLinkDocumentUploadResult(existingFile, link.expiresAt()); + } + + private NormalizedUploadRequest normalize(WorkerLinkDocumentUploadCommand command) { + if (command.size() <= 0) { + throw new ApiException(ErrorCode.VALIDATION_FAILED, "파일 크기는 0보다 커야 합니다."); + } + if (command.size() > MAX_FILE_SIZE_BYTES) { + throw new ApiException(FileErrorCode.FILE_TOO_LARGE); + } + + String fileName = normalizeRequiredText(command.fileName(), "파일명이 필요합니다."); + if (fileName.length() > MAX_FILE_NAME_LENGTH) { + throw new ApiException(ErrorCode.VALIDATION_FAILED, "파일명은 255자 이하여야 합니다."); + } + + String mimeType = normalizeRequiredText(command.mimeType(), "파일 MIME 유형이 필요합니다.") + .toLowerCase(Locale.ROOT); + if (!ALLOWED_MIME_TYPES.contains(mimeType)) { + throw new ApiException(FileErrorCode.UNSUPPORTED_FILE_TYPE); + } + + String purpose = normalizePurpose(command.documentType()); + String idempotencyKey = normalizeRequiredText( + command.idempotencyKey(), + "Idempotency-Key 헤더가 필요합니다." + ); + if (idempotencyKey.length() < MIN_IDEMPOTENCY_KEY_LENGTH + || idempotencyKey.length() > MAX_IDEMPOTENCY_KEY_LENGTH + || !IDEMPOTENCY_KEY_PATTERN.matcher(idempotencyKey).matches()) { + throw new ApiException( + ErrorCode.VALIDATION_FAILED, + "Idempotency-Key는 8~100자의 영문, 숫자, 점, 밑줄, 콜론 또는 하이픈이어야 합니다." + ); + } + + return new NormalizedUploadRequest( + fileName, + mimeType, + command.size(), + purpose, + idempotencyKey + ); + } + + private String normalizePurpose(String documentType) { + if (documentType == null || documentType.isBlank()) { + return DEFAULT_DOCUMENT_PURPOSE; + } + String normalized = documentType.strip().toUpperCase(Locale.ROOT); + try { + return DocumentType.valueOf(normalized).name(); + } catch (IllegalArgumentException exception) { + throw new ApiException(ErrorCode.VALIDATION_FAILED, "지원하지 않는 문서 유형입니다."); + } + } + + private String normalizeRequiredText(String value, String message) { + if (value == null || value.isBlank()) { + throw new ApiException(ErrorCode.VALIDATION_FAILED, message); + } + return value.strip(); + } + + private String storeAndChecksum( + String storageKey, + InputStream content, + NormalizedUploadRequest request, + FileStorageRollbackCompensation.Registration rollbackRegistration + ) { + MessageDigest digest = newSha256Digest(); + CountingInputStream countingContent = new CountingInputStream(content); + try (DigestInputStream digestContent = new DigestInputStream(countingContent, digest)) { + fileStorage.store(storageKey, digestContent, request.size(), request.mimeType()); + rollbackRegistration.markCreated(); + } catch (IOException exception) { + throw new UncheckedIOException("failed to close uploaded file", exception); + } + requireActualSize(countingContent.count(), request.size()); + return HexFormat.of().formatHex(digest.digest()); + } + + private String checksumAndDiscard(InputStream content, long expectedSize) { + MessageDigest digest = newSha256Digest(); + CountingInputStream countingContent = new CountingInputStream(content); + try (DigestInputStream digestContent = new DigestInputStream(countingContent, digest)) { + digestContent.transferTo(OutputStream.nullOutputStream()); + } catch (IOException exception) { + throw new UncheckedIOException("failed to read uploaded file", exception); + } + requireActualSize(countingContent.count(), expectedSize); + return HexFormat.of().formatHex(digest.digest()); + } + + private void requireActualSize(long actualSize, long expectedSize) { + if (actualSize != expectedSize) { + throw new ApiException( + ErrorCode.VALIDATION_FAILED, + "선언된 파일 크기와 실제 파일 크기가 일치하지 않습니다." + ); + } + } + + private String calculateRequestHash(NormalizedUploadRequest request, String contentChecksum) { + MessageDigest digest = newSha256Digest(); + updateDigest(digest, REQUEST_HASH_VERSION); + updateDigest(digest, request.purpose()); + updateDigest(digest, request.fileName()); + updateDigest(digest, request.mimeType()); + updateDigest(digest, Long.toString(request.size())); + updateDigest(digest, contentChecksum); + return HexFormat.of().formatHex(digest.digest()); + } + + private void updateDigest(MessageDigest digest, String value) { + byte[] bytes = value.getBytes(StandardCharsets.UTF_8); + digest.update(ByteBuffer.allocate(Integer.BYTES).putInt(bytes.length).array()); + digest.update(bytes); + } + + private MessageDigest newSha256Digest() { + try { + return MessageDigest.getInstance("SHA-256"); + } catch (NoSuchAlgorithmException exception) { + throw new IllegalStateException("SHA-256 is not available", exception); + } + } + + private record NormalizedUploadRequest( + String fileName, + String mimeType, + long size, + String purpose, + String idempotencyKey + ) { + } + + private static final class CountingInputStream extends FilterInputStream { + + private long count; + + private CountingInputStream(InputStream input) { + super(input); + } + + @Override + public int read() throws IOException { + int value = in.read(); + if (value >= 0) { + count++; + } + return value; + } + + @Override + public int read(byte[] bytes, int offset, int length) throws IOException { + int bytesRead = in.read(bytes, offset, length); + if (bytesRead > 0) { + count += bytesRead; + } + return bytesRead; + } + + private long count() { + return count; + } } } diff --git a/src/main/java/com/fowoco/server/workerlink/application/WorkerLinkDocumentUploadCommand.java b/src/main/java/com/fowoco/server/workerlink/application/WorkerLinkDocumentUploadCommand.java index e87e232d..ae481f92 100644 --- a/src/main/java/com/fowoco/server/workerlink/application/WorkerLinkDocumentUploadCommand.java +++ b/src/main/java/com/fowoco/server/workerlink/application/WorkerLinkDocumentUploadCommand.java @@ -10,6 +10,7 @@ public final class WorkerLinkDocumentUploadCommand { private final long size; private final String documentType; private final String clientRequestId; + private final String idempotencyKey; private final InputStream content; public WorkerLinkDocumentUploadCommand( @@ -19,6 +20,7 @@ public WorkerLinkDocumentUploadCommand( long size, String documentType, String clientRequestId, + String idempotencyKey, InputStream content ) { this.rawToken = rawToken; @@ -27,6 +29,7 @@ public WorkerLinkDocumentUploadCommand( this.size = size; this.documentType = documentType; this.clientRequestId = clientRequestId; + this.idempotencyKey = idempotencyKey; this.content = content; } @@ -54,6 +57,10 @@ public String clientRequestId() { return clientRequestId; } + public String idempotencyKey() { + return idempotencyKey; + } + public InputStream content() { return content; } diff --git a/src/main/java/com/fowoco/server/workerlink/application/error/WorkerLinkErrorCode.java b/src/main/java/com/fowoco/server/workerlink/application/error/WorkerLinkErrorCode.java index 350090eb..49e11cd6 100644 --- a/src/main/java/com/fowoco/server/workerlink/application/error/WorkerLinkErrorCode.java +++ b/src/main/java/com/fowoco/server/workerlink/application/error/WorkerLinkErrorCode.java @@ -48,6 +48,10 @@ public enum WorkerLinkErrorCode implements ApiErrorCode { HttpStatus.UNPROCESSABLE_CONTENT, "업로드된 파일을 찾을 수 없거나 이미 사용된 파일입니다." ), + IDEMPOTENCY_CONFLICT( + HttpStatus.CONFLICT, + "같은 Idempotency-Key가 다른 문서 업로드 요청에 이미 사용되었습니다." + ), WORKER_SLOT_ANSWER_INVALID( HttpStatus.UNPROCESSABLE_CONTENT, "요청하지 않았거나 허용되지 않은 근로자 답변입니다." diff --git a/src/main/java/com/fowoco/server/workerlink/application/port/WorkerDocumentUploadIdempotencyRecord.java b/src/main/java/com/fowoco/server/workerlink/application/port/WorkerDocumentUploadIdempotencyRecord.java new file mode 100644 index 00000000..6c928061 --- /dev/null +++ b/src/main/java/com/fowoco/server/workerlink/application/port/WorkerDocumentUploadIdempotencyRecord.java @@ -0,0 +1,12 @@ +package com.fowoco.server.workerlink.application.port; + +import java.util.Objects; +import java.util.UUID; + +public record WorkerDocumentUploadIdempotencyRecord(UUID storedFileId, String requestHash) { + + public WorkerDocumentUploadIdempotencyRecord { + Objects.requireNonNull(storedFileId, "storedFileId must not be null"); + Objects.requireNonNull(requestHash, "requestHash must not be null"); + } +} diff --git a/src/main/java/com/fowoco/server/workerlink/application/port/WorkerDocumentUploadIdempotencyRepository.java b/src/main/java/com/fowoco/server/workerlink/application/port/WorkerDocumentUploadIdempotencyRepository.java index 23e91c4c..22dc0de5 100644 --- a/src/main/java/com/fowoco/server/workerlink/application/port/WorkerDocumentUploadIdempotencyRepository.java +++ b/src/main/java/com/fowoco/server/workerlink/application/port/WorkerDocumentUploadIdempotencyRepository.java @@ -5,7 +5,17 @@ public interface WorkerDocumentUploadIdempotencyRepository { - Optional findStoredFileId(UUID workerLinkId, UUID companyId, String clientRequestId); + Optional findByKeyHash( + UUID workerLinkId, + UUID companyId, + String idempotencyKeyHash + ); - void save(UUID workerLinkId, UUID companyId, String clientRequestId, UUID storedFileId); + void save( + UUID workerLinkId, + UUID companyId, + String idempotencyKeyHash, + String requestHash, + UUID storedFileId + ); } diff --git a/src/main/java/com/fowoco/server/workerlink/infrastructure/persistence/JpaWorkerDocumentUploadIdempotencyRepository.java b/src/main/java/com/fowoco/server/workerlink/infrastructure/persistence/JpaWorkerDocumentUploadIdempotencyRepository.java index a1a3b2d2..6cd28333 100644 --- a/src/main/java/com/fowoco/server/workerlink/infrastructure/persistence/JpaWorkerDocumentUploadIdempotencyRepository.java +++ b/src/main/java/com/fowoco/server/workerlink/infrastructure/persistence/JpaWorkerDocumentUploadIdempotencyRepository.java @@ -1,6 +1,7 @@ package com.fowoco.server.workerlink.infrastructure.persistence; import com.fowoco.server.workerlink.application.port.WorkerDocumentUploadIdempotencyRepository; +import com.fowoco.server.workerlink.application.port.WorkerDocumentUploadIdempotencyRecord; import jakarta.persistence.EntityManager; import jakarta.persistence.Query; import java.util.Objects; @@ -11,6 +12,8 @@ @Repository public class JpaWorkerDocumentUploadIdempotencyRepository implements WorkerDocumentUploadIdempotencyRepository { + private static final String CANONICAL_CLIENT_REQUEST_ID_PREFIX = "canonical:"; + private final EntityManager entityManager; public JpaWorkerDocumentUploadIdempotencyRepository(EntityManager entityManager) { @@ -18,47 +21,60 @@ public JpaWorkerDocumentUploadIdempotencyRepository(EntityManager entityManager) } @Override - public Optional findStoredFileId( + public Optional findByKeyHash( UUID workerLinkId, UUID companyId, - String clientRequestId + String idempotencyKeyHash ) { Objects.requireNonNull(workerLinkId, "workerLinkId must not be null"); Objects.requireNonNull(companyId, "companyId must not be null"); - Objects.requireNonNull(clientRequestId, "clientRequestId must not be null"); + Objects.requireNonNull(idempotencyKeyHash, "idempotencyKeyHash must not be null"); return entityManager.createNativeQuery( - "SELECT CAST(stored_file_id AS VARCHAR) FROM worker_document_upload_idempotency " + "SELECT CAST(stored_file_id AS VARCHAR), request_hash " + + "FROM worker_document_upload_idempotency " + "WHERE worker_link_id = ?1 AND company_id = ?2 " - + "AND client_request_id = ?3" + + "AND idempotency_key_hash = ?3" ) .setParameter(1, workerLinkId) .setParameter(2, companyId) - .setParameter(3, clientRequestId) + .setParameter(3, idempotencyKeyHash) .getResultStream() .findFirst() - .map(result -> UUID.fromString(result.toString())); + .map(result -> { + Object[] columns = (Object[]) result; + return new WorkerDocumentUploadIdempotencyRecord( + UUID.fromString(columns[0].toString()), + columns[1].toString() + ); + }); } @Override public void save( UUID workerLinkId, UUID companyId, - String clientRequestId, + String idempotencyKeyHash, + String requestHash, UUID storedFileId ) { Objects.requireNonNull(workerLinkId, "workerLinkId must not be null"); Objects.requireNonNull(companyId, "companyId must not be null"); - Objects.requireNonNull(clientRequestId, "clientRequestId must not be null"); + Objects.requireNonNull(idempotencyKeyHash, "idempotencyKeyHash must not be null"); + Objects.requireNonNull(requestHash, "requestHash must not be null"); Objects.requireNonNull(storedFileId, "storedFileId must not be null"); + String compatibilityClientRequestId = CANONICAL_CLIENT_REQUEST_ID_PREFIX + storedFileId; Query query = entityManager.createNativeQuery( "INSERT INTO worker_document_upload_idempotency " - + "(worker_link_id, company_id, client_request_id, stored_file_id) " - + "VALUES (?1, ?2, ?3, ?4)" + + "(worker_link_id, company_id, client_request_id, stored_file_id, " + + "idempotency_key_hash, request_hash) " + + "VALUES (?1, ?2, ?3, ?4, ?5, ?6)" ); query.setParameter(1, workerLinkId); query.setParameter(2, companyId); - query.setParameter(3, clientRequestId); + query.setParameter(3, compatibilityClientRequestId); query.setParameter(4, storedFileId); + query.setParameter(5, idempotencyKeyHash); + query.setParameter(6, requestHash); query.executeUpdate(); } } diff --git a/src/main/resources/db/migration/V51__add_worker_document_upload_idempotency_hashes.sql b/src/main/resources/db/migration/V51__add_worker_document_upload_idempotency_hashes.sql new file mode 100644 index 00000000..1a051265 --- /dev/null +++ b/src/main/resources/db/migration/V51__add_worker_document_upload_idempotency_hashes.sql @@ -0,0 +1,32 @@ +ALTER TABLE worker_document_upload_idempotency + ADD COLUMN idempotency_key_hash VARCHAR(64); + +ALTER TABLE worker_document_upload_idempotency + ADD COLUMN request_hash VARCHAR(64); + +-- Existing rows and the previous application version keep using client_request_id only. +-- New rows written through the canonical Idempotency-Key flow populate both hashes. +ALTER TABLE worker_document_upload_idempotency + ADD CONSTRAINT uq_worker_document_upload_idempotency_key_hash + UNIQUE (worker_link_id, idempotency_key_hash); + +ALTER TABLE worker_document_upload_idempotency + ADD CONSTRAINT ck_worker_document_upload_idempotency_key_hash + CHECK ( + idempotency_key_hash IS NULL + OR CHAR_LENGTH(idempotency_key_hash) = 64 + ); + +ALTER TABLE worker_document_upload_idempotency + ADD CONSTRAINT ck_worker_document_upload_idempotency_request_hash + CHECK ( + request_hash IS NULL + OR CHAR_LENGTH(request_hash) = 64 + ); + +ALTER TABLE worker_document_upload_idempotency + ADD CONSTRAINT ck_worker_document_upload_idempotency_hash_pair + CHECK ( + (idempotency_key_hash IS NULL AND request_hash IS NULL) + OR (idempotency_key_hash IS NOT NULL AND request_hash IS NOT NULL) + ); diff --git a/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java b/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java index 8a4b0030..a7ae54db 100644 --- a/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java +++ b/src/test/java/com/fowoco/server/PostgreSqlMigrationTests.java @@ -301,8 +301,11 @@ private void assertSchemaContract(Connection connection) throws SQLException { .containsEntry("company_id", new ColumnSpec("uuid", false)); assertThat(columnSpecs(connection, "worker_document_upload_idempotency")) .containsEntry("worker_link_id", new ColumnSpec("uuid", false)) + .containsEntry("client_request_id", new ColumnSpec("varchar", false)) .containsEntry("stored_file_id", new ColumnSpec("uuid", false)) - .containsEntry("company_id", new ColumnSpec("uuid", false)); + .containsEntry("company_id", new ColumnSpec("uuid", false)) + .containsEntry("idempotency_key_hash", new ColumnSpec("varchar", true)) + .containsEntry("request_hash", new ColumnSpec("varchar", true)); assertThat(columnSpecs(connection, "workflow_case")) .containsEntry("case_id", new ColumnSpec("uuid", false)) .containsEntry("company_id", new ColumnSpec("uuid", false)) @@ -424,6 +427,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", + "uq_worker_document_upload_idempotency_key_hash", + "ck_worker_document_upload_idempotency_key_hash", + "ck_worker_document_upload_idempotency_request_hash", + "ck_worker_document_upload_idempotency_hash_pair", "pk_outbox_manual_retry", "uq_outbox_manual_retry_event_key", "fk_outbox_manual_retry_event_company", @@ -674,6 +681,54 @@ INSERT INTO worker_link ( TASK_A, COMPANY_A, REVOKED_WORKER_LINK_TOKEN_HASH, USER_A, TASK_A, COMPANY_A, EXPIRED_WORKER_LINK_TOKEN_HASH, USER_A )); + execute(connection, """ + INSERT INTO stored_file ( + stored_file_id, company_id, name, mime_type, size, purpose, + storage_key, scan_status + ) VALUES ( + '23000000-0000-0000-0000-000000000001', '%s', + 'worker-upload.pdf', 'application/pdf', 1, 'WORKER_LINK_SUBMISSION', + 'migration-worker-upload', 'NOT_SCANNED' + ) + """.formatted(COMPANY_A)); + execute(connection, """ + INSERT INTO worker_document_upload_idempotency ( + worker_link_id, company_id, client_request_id, stored_file_id + ) VALUES ( + '21000000-0000-0000-0000-000000000001', '%s', + 'legacy-client-request', '23000000-0000-0000-0000-000000000001' + ) + """.formatted(COMPANY_A)); + execute(connection, """ + INSERT INTO worker_document_upload_idempotency ( + worker_link_id, company_id, client_request_id, stored_file_id, + idempotency_key_hash, request_hash + ) VALUES ( + '21000000-0000-0000-0000-000000000001', '%s', + 'canonical-client-request-a', '23000000-0000-0000-0000-000000000001', + '%s', '%s' + ) + """.formatted(COMPANY_A, "1".repeat(64), "2".repeat(64))); + assertSqlState(connection, "23505", """ + INSERT INTO worker_document_upload_idempotency ( + worker_link_id, company_id, client_request_id, stored_file_id, + idempotency_key_hash, request_hash + ) VALUES ( + '21000000-0000-0000-0000-000000000001', '%s', + 'canonical-client-request-b', '23000000-0000-0000-0000-000000000001', + '%s', '%s' + ) + """.formatted(COMPANY_A, "1".repeat(64), "3".repeat(64))); + assertSqlState(connection, "23514", """ + INSERT INTO worker_document_upload_idempotency ( + worker_link_id, company_id, client_request_id, stored_file_id, + idempotency_key_hash + ) VALUES ( + '21000000-0000-0000-0000-000000000001', '%s', + 'canonical-client-request-c', '23000000-0000-0000-0000-000000000001', + '%s' + ) + """.formatted(COMPANY_A, "4".repeat(64))); assertThat(queryNullableString( connection, "SELECT delivery_status FROM worker_link WHERE worker_link_id = ?::uuid", diff --git a/src/test/java/com/fowoco/server/document/DocumentOcrApiIntegrationTest.java b/src/test/java/com/fowoco/server/document/DocumentOcrApiIntegrationTest.java index a0840806..2d8a7faa 100644 --- a/src/test/java/com/fowoco/server/document/DocumentOcrApiIntegrationTest.java +++ b/src/test/java/com/fowoco/server/document/DocumentOcrApiIntegrationTest.java @@ -529,6 +529,11 @@ public Optional open(String storageKey) { byte[] value = files.get(storageKey); return value == null ? Optional.empty() : Optional.of(new ByteArrayInputStream(value)); } + + @Override + public void deleteIfExists(String storageKey) { + files.remove(storageKey); + } } static final class TestAiOcrClient implements AiOcrClient { diff --git a/src/test/java/com/fowoco/server/file/FileRollbackIntegrationTest.java b/src/test/java/com/fowoco/server/file/FileRollbackIntegrationTest.java new file mode 100644 index 00000000..8d7578d0 --- /dev/null +++ b/src/test/java/com/fowoco/server/file/FileRollbackIntegrationTest.java @@ -0,0 +1,186 @@ +package com.fowoco.server.file; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.reset; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.fowoco.server.audit.application.port.AuditEventRepository; +import com.fowoco.server.auth.application.ActorContext; +import com.fowoco.server.auth.domain.UserRole; +import com.fowoco.server.common.id.UuidGenerator; +import com.fowoco.server.common.web.RequestMetadata; +import com.fowoco.server.file.application.FileCreateCommand; +import com.fowoco.server.file.application.FileService; +import com.fowoco.server.file.application.port.FileStorage; +import com.fowoco.server.file.application.port.StoredFileRepository; +import com.fowoco.server.file.domain.StoredFile; +import java.io.ByteArrayInputStream; +import java.io.InputStream; +import java.io.UncheckedIOException; +import java.nio.file.FileAlreadyExistsException; +import java.nio.charset.StandardCharsets; +import java.util.HashMap; +import java.util.Map; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.atomic.AtomicReference; +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.test.context.ActiveProfiles; +import org.springframework.test.context.bean.override.mockito.MockitoBean; +import org.springframework.transaction.support.TransactionTemplate; + +@ActiveProfiles("test") +@SpringBootTest +class FileRollbackIntegrationTest { + + private static final UUID COMPANY_ID = + UUID.fromString("74000000-0000-0000-0000-000000000001"); + private static final UUID ACTOR_ID = + UUID.fromString("74100000-0000-0000-0000-000000000001"); + private static final ActorContext ACTOR = + new ActorContext(ACTOR_ID, COMPANY_ID, Set.of(UserRole.HR)); + private static final RequestMetadata METADATA = + new RequestMetadata("file-rollback-request", null); + + @Autowired + private FileService fileService; + + @Autowired + private TransactionTemplate transactionTemplate; + + @MockitoBean + private FileStorage fileStorage; + + @MockitoBean + private StoredFileRepository storedFileRepository; + + @MockitoBean + private AuditEventRepository auditRepository; + + @MockitoBean + private UuidGenerator uuidGenerator; + + private final Map storedContents = new HashMap<>(); + + @BeforeEach + void setUpStorage() throws Exception { + reset(fileStorage, storedFileRepository, auditRepository, uuidGenerator); + storedContents.clear(); + when(uuidGenerator.generate()).thenAnswer(invocation -> UUID.randomUUID()); + doAnswer(invocation -> { + String storageKey = invocation.getArgument(0); + InputStream content = invocation.getArgument(1); + storedContents.put(storageKey, content.readAllBytes()); + return null; + }).when(fileStorage).store(anyString(), any(InputStream.class), anyLong(), anyString()); + doAnswer(invocation -> { + storedContents.remove(invocation.getArgument(0)); + return null; + }).when(fileStorage).deleteIfExists(anyString()); + } + + @Test + void keepsPreexistingFileWhenGeneratedStorageKeyCollides() { + UUID existingStoredFileId = UUID.fromString("74200000-0000-0000-0000-000000000001"); + String storageKey = existingStoredFileId.toString(); + byte[] existingContent = "previously-committed-content".getBytes(StandardCharsets.UTF_8); + storedContents.put(storageKey, existingContent); + when(uuidGenerator.generate()).thenReturn(existingStoredFileId); + doThrow(new UncheckedIOException( + "failed to store file: " + storageKey, + new FileAlreadyExistsException(storageKey) + )).when(fileStorage).store(eq(storageKey), any(InputStream.class), anyLong(), anyString()); + + assertThatThrownBy(() -> fileService.upload(command(), ACTOR, METADATA)) + .isInstanceOf(UncheckedIOException.class); + + assertThat(storedContents).containsEntry(storageKey, existingContent); + verify(fileStorage, never()).deleteIfExists(storageKey); + verify(storedFileRepository, never()).insert(any()); + } + + @Test + void keepsStoredFileAfterCommit() { + StoredFile storedFile = fileService.upload(command(), ACTOR, METADATA); + + assertThat(storedContents).containsKey(storedFile.storageKey()); + verify(fileStorage, never()).deleteIfExists(storedFile.storageKey()); + } + + @Test + void cleansStoredFileWhenStoredFileInsertFails() { + RuntimeException originalFailure = new IllegalStateException("stored file insert failed"); + doThrow(originalFailure).when(storedFileRepository).insert(any()); + + assertThatThrownBy(() -> fileService.upload(command(), ACTOR, METADATA)) + .isSameAs(originalFailure); + + assertThat(storedContents).isEmpty(); + } + + @Test + void cleansStoredFileWhenAuditAppendFails() { + RuntimeException originalFailure = new IllegalStateException("audit append failed"); + doThrow(originalFailure).when(auditRepository).append(any()); + + assertThatThrownBy(() -> fileService.upload(command(), ACTOR, METADATA)) + .isSameAs(originalFailure); + + assertThat(storedContents).isEmpty(); + } + + @Test + void cleansStoredFileWhenOuterTransactionFailsAfterUploadReturns() { + RuntimeException originalFailure = new IllegalStateException("outer transaction failed"); + AtomicReference storageKey = new AtomicReference<>(); + + assertThatThrownBy(() -> transactionTemplate.executeWithoutResult(status -> { + StoredFile storedFile = fileService.upload(command(), ACTOR, METADATA); + storageKey.set(storedFile.storageKey()); + assertThat(storedContents).containsKey(storedFile.storageKey()); + throw originalFailure; + })).isSameAs(originalFailure); + + assertThat(storageKey.get()).isNotNull(); + assertThat(storedContents).doesNotContainKey(storageKey.get()); + } + + @Test + void cleanupFailureDoesNotReplaceOriginalTransactionFailure() { + RuntimeException originalFailure = new IllegalStateException("stored file insert failed"); + RuntimeException cleanupFailure = new IllegalStateException("file cleanup failed"); + doThrow(originalFailure).when(storedFileRepository).insert(any()); + doThrow(cleanupFailure).when(fileStorage).deleteIfExists(anyString()); + + assertThatThrownBy(() -> fileService.upload(command(), ACTOR, METADATA)) + .isSameAs(originalFailure); + + assertThat(storedContents).hasSize(1); + } + + private FileCreateCommand command() { + byte[] content = "rollback-safe-content".getBytes(StandardCharsets.UTF_8); + return new FileCreateCommand( + COMPANY_ID, + "rollback-test.pdf", + "application/pdf", + content.length, + "GENERAL", + null, + null, + new ByteArrayInputStream(content) + ); + } +} diff --git a/src/test/java/com/fowoco/server/file/application/FileStorageRollbackCompensationTest.java b/src/test/java/com/fowoco/server/file/application/FileStorageRollbackCompensationTest.java new file mode 100644 index 00000000..adaf072d --- /dev/null +++ b/src/test/java/com/fowoco/server/file/application/FileStorageRollbackCompensationTest.java @@ -0,0 +1,111 @@ +package com.fowoco.server.file.application; + +import static org.assertj.core.api.Assertions.assertThatCode; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; + +import com.fowoco.server.common.web.RequestMetadata; +import com.fowoco.server.file.application.port.FileStorage; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.transaction.support.TransactionSynchronization; +import org.springframework.transaction.support.TransactionSynchronizationManager; + +class FileStorageRollbackCompensationTest { + + private static final String STORAGE_KEY = "server-generated-storage-key"; + private static final RequestMetadata METADATA = new RequestMetadata("request-1", null); + private static final String ACTION = "file_upload"; + + private FileStorage fileStorage; + private FileStorageRollbackCompensation compensation; + + @BeforeEach + void setUp() { + TransactionSynchronizationManager.clear(); + TransactionSynchronizationManager.setActualTransactionActive(true); + TransactionSynchronizationManager.initSynchronization(); + fileStorage = mock(FileStorage.class); + compensation = new FileStorageRollbackCompensation(fileStorage); + } + + @AfterEach + void tearDown() { + TransactionSynchronizationManager.clear(); + } + + @Test + void keepsFileAfterCommit() { + FileStorageRollbackCompensation.Registration registration = + compensation.register(STORAGE_KEY, METADATA, ACTION); + registration.markCreated(); + + synchronization().afterCompletion(TransactionSynchronization.STATUS_COMMITTED); + + verify(fileStorage, never()).deleteIfExists(STORAGE_KEY); + } + + @Test + void deletesFileAfterRollback() { + FileStorageRollbackCompensation.Registration registration = + compensation.register(STORAGE_KEY, METADATA, ACTION); + registration.markCreated(); + + synchronization().afterCompletion(TransactionSynchronization.STATUS_ROLLED_BACK); + + verify(fileStorage).deleteIfExists(STORAGE_KEY); + } + + @Test + void keepsFileWhenStoreDidNotCreateIt() { + compensation.register(STORAGE_KEY, METADATA, ACTION); + + synchronization().afterCompletion(TransactionSynchronization.STATUS_ROLLED_BACK); + + verify(fileStorage, never()).deleteIfExists(STORAGE_KEY); + } + + @Test + void keepsFileWhenTransactionCompletionIsUnknown() { + FileStorageRollbackCompensation.Registration registration = + compensation.register(STORAGE_KEY, METADATA, ACTION); + registration.markCreated(); + + synchronization().afterCompletion(TransactionSynchronization.STATUS_UNKNOWN); + + verify(fileStorage, never()).deleteIfExists(STORAGE_KEY); + } + + @Test + void cleanupFailureDoesNotEscapeTransactionCompletionCallback() { + RuntimeException cleanupFailure = new IllegalStateException("cleanup unavailable"); + doThrow(cleanupFailure).when(fileStorage).deleteIfExists(STORAGE_KEY); + FileStorageRollbackCompensation.Registration registration = + compensation.register(STORAGE_KEY, METADATA, ACTION); + registration.markCreated(); + + assertThatCode(() -> synchronization().afterCompletion(TransactionSynchronization.STATUS_ROLLED_BACK)) + .doesNotThrowAnyException(); + + verify(fileStorage).deleteIfExists(STORAGE_KEY); + } + + @Test + void rejectsRegistrationWithoutActiveTransactionSynchronization() { + TransactionSynchronizationManager.clear(); + + assertThatThrownBy(() -> compensation.register(STORAGE_KEY, METADATA, ACTION)) + .isInstanceOf(IllegalStateException.class) + .hasMessageContaining("active transaction synchronization"); + + verify(fileStorage, never()).deleteIfExists(STORAGE_KEY); + } + + private TransactionSynchronization synchronization() { + return TransactionSynchronizationManager.getSynchronizations().get(0); + } +} diff --git a/src/test/java/com/fowoco/server/file/infrastructure/LocalFileStorageTest.java b/src/test/java/com/fowoco/server/file/infrastructure/LocalFileStorageTest.java new file mode 100644 index 00000000..72adcb52 --- /dev/null +++ b/src/test/java/com/fowoco/server/file/infrastructure/LocalFileStorageTest.java @@ -0,0 +1,162 @@ +package com.fowoco.server.file.infrastructure; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import java.io.ByteArrayInputStream; +import java.io.IOException; +import java.io.InputStream; +import java.io.UncheckedIOException; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.List; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.io.TempDir; + +class LocalFileStorageTest { + + private static final String MIME_TYPE = "application/pdf"; + + @TempDir + Path storageRoot; + + @Test + void storesCompleteContentWithoutTemporaryArtifacts() throws Exception { + LocalFileStorage storage = storage(); + byte[] content = "complete-content".getBytes(StandardCharsets.UTF_8); + + storage.store("stored-file-id", new ByteArrayInputStream(content), content.length, MIME_TYPE); + + assertThat(Files.readAllBytes(storageRoot.resolve("stored-file-id"))).isEqualTo(content); + assertThat(temporaryArtifacts()).isEmpty(); + } + + @Test + void cleansTemporaryArtifactWhenSourceReadFails() throws Exception { + LocalFileStorage storage = storage(); + + assertThatThrownBy(() -> storage.store("stored-file-id", failingInput(), 8, MIME_TYPE)) + .isInstanceOf(UncheckedIOException.class) + .hasMessageContaining("failed to store file"); + + assertThat(Files.exists(storageRoot.resolve("stored-file-id"))).isFalse(); + assertThat(temporaryArtifacts()).isEmpty(); + } + + @Test + void cleansTemporaryArtifactWhenStoredSizeDiffersFromDeclaredSize() throws Exception { + LocalFileStorage storage = storage(); + byte[] content = "content".getBytes(StandardCharsets.UTF_8); + + assertThatThrownBy(() -> storage.store( + "stored-file-id", + new ByteArrayInputStream(content), + content.length + 1, + MIME_TYPE + )) + .isInstanceOf(UncheckedIOException.class) + .hasMessageContaining("failed to store file"); + + assertThat(Files.exists(storageRoot.resolve("stored-file-id"))).isFalse(); + assertThat(temporaryArtifacts()).isEmpty(); + } + + @Test + void cleansTemporaryArtifactWhenFinalPathCannotBeCreated() throws Exception { + LocalFileStorage storage = storage(); + byte[] content = "content".getBytes(StandardCharsets.UTF_8); + + assertThatThrownBy(() -> storage.store( + "missing-directory/stored-file-id", + new ByteArrayInputStream(content), + content.length, + MIME_TYPE + )) + .isInstanceOf(UncheckedIOException.class) + .hasMessageContaining("failed to store file"); + + assertThat(Files.exists(storageRoot.resolve("missing-directory/stored-file-id"))).isFalse(); + assertThat(temporaryArtifacts()).isEmpty(); + } + + @Test + void doesNotReplaceAlreadyExistingFinalFile() throws Exception { + LocalFileStorage storage = storage(); + byte[] existing = "existing".getBytes(StandardCharsets.UTF_8); + byte[] replacement = "replacement".getBytes(StandardCharsets.UTF_8); + Files.write(storageRoot.resolve("stored-file-id"), existing); + + assertThatThrownBy(() -> storage.store( + "stored-file-id", + new ByteArrayInputStream(replacement), + replacement.length, + MIME_TYPE + )).isInstanceOf(UncheckedIOException.class); + + assertThat(Files.readAllBytes(storageRoot.resolve("stored-file-id"))).isEqualTo(existing); + assertThat(temporaryArtifacts()).isEmpty(); + } + + @Test + void deletesOnlyRequestedFileAndTreatsRepeatedDeletionAsSuccess() throws Exception { + LocalFileStorage storage = storage(); + byte[] content = "content".getBytes(StandardCharsets.UTF_8); + storage.store("owned-file", new ByteArrayInputStream(content), content.length, MIME_TYPE); + storage.store("other-file", new ByteArrayInputStream(content), content.length, MIME_TYPE); + + storage.deleteIfExists("owned-file"); + storage.deleteIfExists("owned-file"); + + assertThat(Files.exists(storageRoot.resolve("owned-file"))).isFalse(); + assertThat(Files.readAllBytes(storageRoot.resolve("other-file"))).isEqualTo(content); + } + + @Test + void rejectsStorageKeysThatEscapeTheRoot() { + LocalFileStorage storage = storage(); + + assertThatThrownBy(() -> storage.deleteIfExists("../outside-file")) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("must not escape"); + } + + private LocalFileStorage storage() { + return new LocalFileStorage(storageRoot.toString()); + } + + private List temporaryArtifacts() throws IOException { + try (var files = Files.list(storageRoot)) { + return files + .filter(path -> path.getFileName().toString().startsWith(".fowoco-upload-")) + .toList(); + } + } + + private InputStream failingInput() { + return new InputStream() { + private int reads; + + @Override + public int read() throws IOException { + if (reads++ >= 3) { + throw new IOException("simulated source failure"); + } + return 'a'; + } + + @Override + public int read(byte[] buffer, int offset, int length) throws IOException { + if (reads >= 3) { + throw new IOException("simulated source failure"); + } + int written = Math.min(length, 3 - reads); + for (int index = 0; index < written; index++) { + buffer[offset + index] = 'a'; + } + reads += written; + return written; + } + }; + } +} diff --git a/src/test/java/com/fowoco/server/file/support/FakeFileStorage.java b/src/test/java/com/fowoco/server/file/support/FakeFileStorage.java index c7cd95cc..1ba80b4e 100644 --- a/src/test/java/com/fowoco/server/file/support/FakeFileStorage.java +++ b/src/test/java/com/fowoco/server/file/support/FakeFileStorage.java @@ -29,6 +29,11 @@ public Optional open(String storageKey) { return content == null ? Optional.empty() : Optional.of(new java.io.ByteArrayInputStream(content)); } + @Override + public void deleteIfExists(String storageKey) { + storedContents.remove(storageKey); + } + public boolean contains(String storageKey) { return storedContents.containsKey(storageKey); } diff --git a/src/test/java/com/fowoco/server/workerlink/WorkerLinkDocumentPostgreSqlIntegrationTest.java b/src/test/java/com/fowoco/server/workerlink/WorkerLinkDocumentPostgreSqlIntegrationTest.java new file mode 100644 index 00000000..1c491c1a --- /dev/null +++ b/src/test/java/com/fowoco/server/workerlink/WorkerLinkDocumentPostgreSqlIntegrationTest.java @@ -0,0 +1,535 @@ +package com.fowoco.server.workerlink; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +import com.fowoco.server.auth.application.ActorContext; +import com.fowoco.server.auth.domain.UserRole; +import com.fowoco.server.common.web.RequestMetadata; +import com.fowoco.server.file.application.FileCreateCommand; +import com.fowoco.server.file.application.FileService; +import com.fowoco.server.file.domain.StoredFile; +import com.fowoco.server.workerlink.application.port.WorkerLinkRepository; +import com.fowoco.server.workerlink.domain.WorkerLink; +import com.fowoco.server.workerlink.infrastructure.persistence.JpaWorkerLinkRepository; +import com.fowoco.server.workerlink.infrastructure.security.WorkerLinkHasher; +import com.jayway.jsonpath.JsonPath; +import java.io.ByteArrayInputStream; +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.HttpRequest; +import java.net.http.HttpResponse; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.util.List; +import java.util.Optional; +import java.util.Set; +import java.util.UUID; +import java.util.concurrent.BrokenBarrierException; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicReference; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.condition.EnabledIfEnvironmentVariable; +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.context.annotation.Primary; +import org.springframework.http.HttpHeaders; +import org.springframework.jdbc.core.JdbcTemplate; +import org.springframework.test.context.ActiveProfiles; +import org.springframework.test.context.DynamicPropertyRegistry; +import org.springframework.test.context.DynamicPropertySource; +import org.springframework.transaction.support.TransactionTemplate; + +@EnabledIfEnvironmentVariable(named = "POSTGRES_TEST_URL", matches = ".+") +@EnabledIfEnvironmentVariable(named = "POSTGRES_TEST_USERNAME", matches = ".+") +@EnabledIfEnvironmentVariable(named = "POSTGRES_TEST_PASSWORD", matches = ".+") +@ActiveProfiles("test") +@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.RANDOM_PORT) +@Import(WorkerLinkDocumentPostgreSqlIntegrationTest.ConcurrencyConfiguration.class) +class WorkerLinkDocumentPostgreSqlIntegrationTest { + + private static final UUID COMPANY_ID = UUID.fromString("17200000-0000-0000-0000-000000000001"); + private static final UUID USER_ID = UUID.fromString("17200000-0000-0000-0000-000000000002"); + private static final UUID WORKER_ID = UUID.fromString("17200000-0000-0000-0000-000000000003"); + private static final UUID CASE_ID = UUID.fromString("17200000-0000-0000-0000-000000000004"); + private static final UUID TASK_ID = UUID.fromString("17200000-0000-0000-0000-000000000005"); + private static final UUID WORKER_LINK_ID = UUID.fromString("17200000-0000-0000-0000-000000000006"); + private static final String RAW_WORKER_LINK_TOKEN = "worker-link-postgresql-concurrency-token"; + private static final String BOUNDARY = "FowocoPostgreSqlUploadBoundary172"; + private static final int CONCURRENCY_ATTEMPTS = 5; + private static final Path STORAGE_ROOT = Path.of( + "build", "test-file-storage", "worker-link-postgresql-" + UUID.randomUUID() + ).toAbsolutePath().normalize(); + + @LocalServerPort + private int port; + + @Autowired + private JdbcTemplate jdbcTemplate; + + @Autowired + private WorkerLinkHasher workerLinkHasher; + + @Autowired + private CoordinatedWorkerLinkRepository coordinatedWorkerLinkRepository; + + @Autowired + private FileService fileService; + + @Autowired + private TransactionTemplate transactionTemplate; + + private final HttpClient httpClient = HttpClient.newHttpClient(); + + @DynamicPropertySource + static void usePostgreSqlAndIsolatedFileStorage(DynamicPropertyRegistry registry) { + registry.add("spring.datasource.url", () -> requiredEnvironmentVariable("POSTGRES_TEST_URL")); + registry.add("spring.datasource.username", () -> requiredEnvironmentVariable("POSTGRES_TEST_USERNAME")); + registry.add("spring.datasource.password", () -> requiredEnvironmentVariable("POSTGRES_TEST_PASSWORD")); + registry.add("spring.datasource.driver-class-name", () -> "org.postgresql.Driver"); + registry.add( + "spring.flyway.locations", + () -> "classpath:db/migration,classpath:db/migration-postgresql" + ); + registry.add("app.file-storage.local-path", STORAGE_ROOT::toString); + } + + @BeforeEach + void seedWorkerLink() throws IOException { + assertPostgreSql16(); + clearDatabaseRows(); + clearStorageDirectory(); + + jdbcTemplate.update( + """ + INSERT INTO company (company_id, name, status, created_at, updated_at, version) + VALUES (?, '파일 보상 동시성 테스트 사업장', 'ACTIVE', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, 0) + """, + COMPANY_ID + ); + jdbcTemplate.update( + """ + INSERT INTO user_account ( + user_id, company_id, email, normalized_email, password_hash, + role, status, created_at, updated_at, version + ) VALUES ( + ?, ?, 'file.rollback.172@example.com', 'file.rollback.172@example.com', + 'unused-test-password-hash', 'HR', 'ACTIVE', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, 0 + ) + """, + USER_ID, + COMPANY_ID + ); + jdbcTemplate.update( + """ + INSERT INTO worker ( + worker_id, company_id, display_name, work_status, created_at, updated_at, version + ) VALUES (?, ?, 'PostgreSQL 동시성 테스트 근로자', 'ACTIVE', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, 0) + """, + WORKER_ID, + COMPANY_ID + ); + jdbcTemplate.update( + """ + INSERT INTO workflow_case ( + case_id, company_id, worker_id, title, lifecycle_status, + priority, workflow_catalog_version, workflow_snapshot_json, + created_by, created_at, updated_at, version + ) VALUES ( + ?, ?, ?, 'Worker Link 문서 제출', 'ACTIVE', 'NORMAL', + '2026.07', '{}', ?, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, 0 + ) + """, + CASE_ID, + COMPANY_ID, + WORKER_ID, + USER_ID + ); + jdbcTemplate.update( + """ + INSERT INTO task ( + task_id, company_id, worker_id, case_id, task_type, + workflow_id, workflow_catalog_version, title, business_data_json, + critical_fingerprint, content_revision, source, status, + created_by, updated_by, created_at, updated_at, version + ) VALUES ( + ?, ?, ?, ?, 'RECONTRACT', 'WF-CON-001', '2026.07', + 'Worker Link 문서 제출', '{}', ?, 0, 'MANUAL', 'APPROVED', + ?, ?, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, 0 + ) + """, + TASK_ID, + COMPANY_ID, + WORKER_ID, + CASE_ID, + "f".repeat(64), + USER_ID, + USER_ID + ); + jdbcTemplate.update( + """ + INSERT INTO worker_link ( + worker_link_id, task_id, company_id, token_hash, expires_at, + status, conversation_status, issued_by, idempotency_key, + created_at, updated_at, version + ) VALUES ( + ?, ?, ?, ?, CURRENT_TIMESTAMP + INTERVAL '1 day', 'ACTIVE', + 'WAITING_WORKER', ?, 'postgresql-concurrency-link', + CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, 0 + ) + """, + WORKER_LINK_ID, + TASK_ID, + COMPANY_ID, + workerLinkHasher.hash(RAW_WORKER_LINK_TOKEN), + USER_ID + ); + } + + @AfterEach + void cleanUp() throws IOException { + clearDatabaseRows(); + clearStorageDirectory(); + } + + @Test + void concurrentSameKeyUploadsConvergeOnOnePostgreSqlRowAndOneLocalFile() throws Exception { + String lastIdempotencyKey = null; + byte[] lastContent = null; + + for (int attempt = 1; attempt <= CONCURRENCY_ATTEMPTS; attempt++) { + String idempotencyKey = "postgresql-upload-key-" + attempt; + byte[] content = ("concurrent-upload-content-" + attempt).getBytes(StandardCharsets.UTF_8); + + coordinatedWorkerLinkRepository.coordinateNextTwoTokenLookups(); + List> responses = concurrentlyUpload( + idempotencyKey, + "passport-" + attempt + ".pdf", + content + ); + + assertThat(responses).extracting(HttpResponse::statusCode).containsExactly(201, 201); + String firstUploadId = JsonPath.read(responses.get(0).body(), "$.upload_id"); + String secondUploadId = JsonPath.read(responses.get(1).body(), "$.upload_id"); + assertThat(secondUploadId).isEqualTo(firstUploadId); + + assertThat(idempotencyRowCount()).isEqualTo(attempt); + assertThat(storedFileRowCount()).isEqualTo(attempt); + assertThat(finalFileCount()).isEqualTo(attempt); + assertThat(temporaryFileCount()).isZero(); + + String storageKey = jdbcTemplate.queryForObject( + "SELECT storage_key FROM stored_file WHERE stored_file_id = ?", + String.class, + UUID.fromString(firstUploadId) + ); + assertThat(Files.readAllBytes(STORAGE_ROOT.resolve(storageKey))).isEqualTo(content); + + lastIdempotencyKey = idempotencyKey; + lastContent = content; + } + + HttpResponse conflict = upload( + lastIdempotencyKey, + "passport-" + CONCURRENCY_ATTEMPTS + ".pdf", + differentBytesWithSameLength(lastContent) + ); + + assertThat(conflict.statusCode()).as(conflict.body()).isEqualTo(409); + assertThat(JsonPath.read(conflict.body(), "$.code")).isEqualTo("IDEMPOTENCY_CONFLICT"); + assertThat(idempotencyRowCount()).isEqualTo(CONCURRENCY_ATTEMPTS); + assertThat(storedFileRowCount()).isEqualTo(CONCURRENCY_ATTEMPTS); + assertThat(finalFileCount()).isEqualTo(CONCURRENCY_ATTEMPTS); + assertThat(temporaryFileCount()).isZero(); + } + + @Test + void transactionRollbackDeletesTheActualLocalFileAndDatabaseRows() { + byte[] content = "rollback-file-content".getBytes(StandardCharsets.UTF_8); + AtomicReference uploadedFile = new AtomicReference<>(); + ActorContext actor = new ActorContext(USER_ID, COMPANY_ID, Set.of(UserRole.HR)); + + assertThatThrownBy(() -> transactionTemplate.executeWithoutResult(status -> { + StoredFile storedFile = fileService.upload( + new FileCreateCommand( + COMPANY_ID, + "rollback.pdf", + "application/pdf", + content.length, + "ROLLBACK_INTEGRATION_TEST", + null, + null, + new ByteArrayInputStream(content) + ), + actor, + new RequestMetadata("file-storage-rollback-postgresql", null) + ); + uploadedFile.set(storedFile); + throw new ForcedRollbackException(); + })).isInstanceOf(ForcedRollbackException.class); + + StoredFile rolledBackFile = uploadedFile.get(); + assertThat(rolledBackFile).isNotNull(); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM stored_file WHERE stored_file_id = ?", + Integer.class, + rolledBackFile.storedFileId() + )).isZero(); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM audit_event WHERE request_id = 'file-storage-rollback-postgresql'", + Integer.class + )).isZero(); + assertThat(STORAGE_ROOT.resolve(rolledBackFile.storageKey())).doesNotExist(); + assertThat(finalFileCount()).isZero(); + assertThat(temporaryFileCount()).isZero(); + } + + private List> concurrentlyUpload( + String idempotencyKey, + String filename, + byte[] content + ) throws Exception { + ExecutorService executor = Executors.newFixedThreadPool(2); + CountDownLatch workersReady = new CountDownLatch(2); + CountDownLatch start = new CountDownLatch(1); + try { + Future> first = executor.submit( + () -> uploadAfterSignal(idempotencyKey, filename, content, workersReady, start) + ); + Future> second = executor.submit( + () -> uploadAfterSignal(idempotencyKey, filename, content, workersReady, start) + ); + + assertThat(workersReady.await(5, TimeUnit.SECONDS)) + .as("both uploads must be ready before they are released") + .isTrue(); + start.countDown(); + return List.of(first.get(30, TimeUnit.SECONDS), second.get(30, TimeUnit.SECONDS)); + } finally { + start.countDown(); + executor.shutdownNow(); + } + } + + private HttpResponse uploadAfterSignal( + String idempotencyKey, + String filename, + byte[] content, + CountDownLatch workersReady, + CountDownLatch start + ) throws Exception { + workersReady.countDown(); + if (!start.await(5, TimeUnit.SECONDS)) { + throw new IllegalStateException("upload concurrency start signal timed out"); + } + return upload(idempotencyKey, filename, content); + } + + private HttpResponse upload(String idempotencyKey, String filename, byte[] content) throws Exception { + ByteArrayOutputStream body = new ByteArrayOutputStream(); + body.write(("--" + BOUNDARY + "\r\n").getBytes(StandardCharsets.UTF_8)); + body.write(("Content-Disposition: form-data; name=\"file\"; filename=\"" + + filename + "\"\r\n").getBytes(StandardCharsets.UTF_8)); + body.write("Content-Type: application/pdf\r\n\r\n".getBytes(StandardCharsets.UTF_8)); + body.write(content); + body.write("\r\n".getBytes(StandardCharsets.UTF_8)); + body.write(("--" + BOUNDARY + "--\r\n").getBytes(StandardCharsets.UTF_8)); + + HttpRequest request = HttpRequest.newBuilder( + URI.create("http://localhost:" + port + + "/api/v1/public/worker-links/" + RAW_WORKER_LINK_TOKEN + "/documents") + ) + .header(HttpHeaders.CONTENT_TYPE, "multipart/form-data; boundary=" + BOUNDARY) + .header("Idempotency-Key", idempotencyKey) + .POST(HttpRequest.BodyPublishers.ofByteArray(body.toByteArray())) + .build(); + return httpClient.send(request, HttpResponse.BodyHandlers.ofString()); + } + + private int idempotencyRowCount() { + return jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM worker_document_upload_idempotency WHERE worker_link_id = ?", + Integer.class, + WORKER_LINK_ID + ); + } + + private int storedFileRowCount() { + return jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM stored_file WHERE company_id = ?", + Integer.class, + COMPANY_ID + ); + } + + private long finalFileCount() { + return storagePaths().stream() + .filter(path -> !path.getFileName().toString().startsWith(".fowoco-upload-")) + .count(); + } + + private long temporaryFileCount() { + return storagePaths().stream() + .filter(path -> path.getFileName().toString().startsWith(".fowoco-upload-")) + .count(); + } + + private List storagePaths() { + try (var paths = Files.list(STORAGE_ROOT)) { + return paths.toList(); + } catch (IOException exception) { + throw new IllegalStateException("failed to inspect the integration-test storage directory", exception); + } + } + + private void clearStorageDirectory() throws IOException { + Files.createDirectories(STORAGE_ROOT); + try (var paths = Files.list(STORAGE_ROOT)) { + for (Path path : paths.toList()) { + if (!path.toAbsolutePath().normalize().startsWith(STORAGE_ROOT)) { + throw new IllegalStateException("test storage path escaped its isolated root"); + } + Files.deleteIfExists(path); + } + } + } + + private void clearDatabaseRows() { + jdbcTemplate.update("DELETE FROM worker_document_upload_idempotency WHERE company_id = ?", COMPANY_ID); + jdbcTemplate.update("DELETE FROM audit_event WHERE company_id = ?", COMPANY_ID); + jdbcTemplate.update("DELETE FROM stored_file WHERE company_id = ?", COMPANY_ID); + jdbcTemplate.update("DELETE FROM worker_link WHERE company_id = ?", COMPANY_ID); + jdbcTemplate.update("DELETE FROM task WHERE company_id = ?", COMPANY_ID); + jdbcTemplate.update("DELETE FROM workflow_case WHERE company_id = ?", COMPANY_ID); + jdbcTemplate.update("DELETE FROM worker WHERE company_id = ?", COMPANY_ID); + jdbcTemplate.update("DELETE FROM user_account WHERE company_id = ?", COMPANY_ID); + jdbcTemplate.update("DELETE FROM company WHERE company_id = ?", COMPANY_ID); + } + + private void assertPostgreSql16() { + Integer versionNumber = jdbcTemplate.queryForObject( + "SELECT current_setting('server_version_num')::integer", + Integer.class + ); + assertThat(versionNumber) + .as("this concurrency test must run against PostgreSQL 16") + .isBetween(160000, 169999); + } + + private static byte[] differentBytesWithSameLength(byte[] original) { + byte[] different = original.clone(); + different[0] = different[0] == 'x' ? (byte) 'y' : (byte) 'x'; + return different; + } + + private static String requiredEnvironmentVariable(String name) { + String value = System.getenv(name); + if (value == null || value.isBlank()) { + throw new IllegalStateException(name + " environment variable is required."); + } + return value; + } + + @TestConfiguration(proxyBeanMethods = false) + static class ConcurrencyConfiguration { + + @Bean + @Primary + CoordinatedWorkerLinkRepository coordinatedWorkerLinkRepository(JpaWorkerLinkRepository delegate) { + return new CoordinatedWorkerLinkRepository(delegate); + } + } + + static final class CoordinatedWorkerLinkRepository implements WorkerLinkRepository { + + private final WorkerLinkRepository delegate; + private final AtomicInteger coordinatedLookupCount = new AtomicInteger(); + private volatile CyclicBarrier lookupBarrier; + + CoordinatedWorkerLinkRepository(WorkerLinkRepository delegate) { + this.delegate = delegate; + } + + void coordinateNextTwoTokenLookups() { + coordinatedLookupCount.set(0); + lookupBarrier = new CyclicBarrier(2); + } + + @Override + public void insert(WorkerLink workerLink) { + delegate.insert(workerLink); + } + + @Override + public WorkerLink update(WorkerLink workerLink) { + return delegate.update(workerLink); + } + + @Override + public Optional findByTokenHash(String tokenHash) { + Optional result = delegate.findByTokenHash(tokenHash); + awaitCoordinatedLookupIfRequired(); + return result; + } + + @Override + public Optional findByIdAndCompanyId(UUID workerLinkId, UUID companyId) { + return delegate.findByIdAndCompanyId(workerLinkId, companyId); + } + + @Override + public Optional findByIdAndCompanyIdForUpdate(UUID workerLinkId, UUID companyId) { + return delegate.findByIdAndCompanyIdForUpdate(workerLinkId, companyId); + } + + @Override + public Optional findActiveByTaskIdAndCompanyId(UUID taskId, UUID companyId) { + return delegate.findActiveByTaskIdAndCompanyId(taskId, companyId); + } + + @Override + public Optional findByTaskIdAndIdempotencyKey(UUID taskId, String idempotencyKey) { + return delegate.findByTaskIdAndIdempotencyKey(taskId, idempotencyKey); + } + + @Override + public List findAllByTaskIdAndCompanyId(UUID taskId, UUID companyId) { + return delegate.findAllByTaskIdAndCompanyId(taskId, companyId); + } + + private void awaitCoordinatedLookupIfRequired() { + CyclicBarrier currentBarrier = lookupBarrier; + int lookupNumber = coordinatedLookupCount.incrementAndGet(); + if (currentBarrier == null || lookupNumber > 2) { + return; + } + try { + currentBarrier.await(10, TimeUnit.SECONDS); + } catch (InterruptedException exception) { + Thread.currentThread().interrupt(); + throw new IllegalStateException("worker-link lookup coordination was interrupted", exception); + } catch (BrokenBarrierException | TimeoutException exception) { + throw new IllegalStateException("two worker-link lookups did not overlap", exception); + } + } + } + + private static final class ForcedRollbackException extends RuntimeException { + } +} diff --git a/src/test/java/com/fowoco/server/workerlink/WorkerLinkSecurityIntegrationTest.java b/src/test/java/com/fowoco/server/workerlink/WorkerLinkSecurityIntegrationTest.java index 5f6ad982..c6af7f2a 100644 --- a/src/test/java/com/fowoco/server/workerlink/WorkerLinkSecurityIntegrationTest.java +++ b/src/test/java/com/fowoco/server/workerlink/WorkerLinkSecurityIntegrationTest.java @@ -17,6 +17,7 @@ import com.fowoco.server.workerlink.application.port.WorkerLinkSmsMessage; import com.fowoco.server.workerlink.application.port.WorkerLinkSmsProviderException; import com.fowoco.server.workerlink.application.port.WorkerLinkSmsSender; +import com.fowoco.server.workerlink.infrastructure.security.WorkerLinkHasher; import com.jayway.jsonpath.JsonPath; import java.io.ByteArrayOutputStream; import java.io.IOException; @@ -70,6 +71,9 @@ class WorkerLinkSecurityIntegrationTest { @Autowired private PasswordEncoder passwordEncoder; + @Autowired + private WorkerLinkHasher workerLinkHasher; + @MockitoBean private WorkerLinkSmsSender workerLinkSmsSender; @@ -101,6 +105,7 @@ void resetState() { 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 worker_document"); jdbcTemplate.update("DELETE FROM stored_file"); jdbcTemplate.update("DELETE FROM event_consumption"); jdbcTemplate.update("DELETE FROM event_publication"); @@ -111,7 +116,6 @@ void resetState() { jdbcTemplate.update("DELETE FROM task_transition_history"); jdbcTemplate.update("DELETE FROM task_checklist_item"); jdbcTemplate.update("DELETE FROM task"); - jdbcTemplate.update("DELETE FROM worker_document"); jdbcTemplate.update("DELETE FROM worker"); jdbcTemplate.update("UPDATE company_settings SET link_expiry_hours = 72"); } @@ -167,16 +171,20 @@ void fullFlow_issueViewUploadRespond_succeeds() throws Exception { assertThat(uploadResponse.statusCode()).isEqualTo(201); String uploadId = JsonPath.read(uploadResponse.body(), "$.upload_id"); - HttpResponse uploadFirstResponse = uploadFileWithFixedClientRequestId( + HttpResponse uploadFirstResponse = uploadFileWithIdempotencyKey( rawToken, "passport.pdf", "application/pdf", - "content".getBytes(StandardCharsets.UTF_8), "fixed-client-request-id" + "content".getBytes(StandardCharsets.UTF_8), + "legacy-client-request-id", + "fixed-upload-idempotency-key" ); assertThat(uploadFirstResponse.statusCode()).isEqualTo(201); String firstUploadId = JsonPath.read(uploadFirstResponse.body(), "$.upload_id"); - HttpResponse uploadRetryResponse = uploadFileWithFixedClientRequestId( + HttpResponse uploadRetryResponse = uploadFileWithIdempotencyKey( rawToken, "passport.pdf", "application/pdf", - "content".getBytes(StandardCharsets.UTF_8), "fixed-client-request-id" + "content".getBytes(StandardCharsets.UTF_8), + null, + "fixed-upload-idempotency-key" ); assertThat(uploadRetryResponse.statusCode()).isEqualTo(201); String retryUploadId = JsonPath.read(uploadRetryResponse.body(), "$.upload_id"); @@ -242,6 +250,191 @@ void fullFlow_issueViewUploadRespond_succeeds() throws Exception { assertThat(otherCompanyActivities.statusCode()).isEqualTo(404); } + @Test + void documentUploadUsesCanonicalIdempotencyKeyAndRejectsDifferentReplay() throws Exception { + String hrToken = accessToken(login(HR_A_EMAIL)); + String workerId = registerWorker(hrToken, "업로드멱등성테스트근로자"); + String taskId = createApprovedTask(hrToken, workerId); + saveDocumentRequestDraft(hrToken, taskId); + String rawToken = issueWorkerLink(hrToken, taskId, "document-upload-idempotency-link-key"); + + HttpResponse missingHeader = uploadFileWithIdempotencyKey( + rawToken, + "passport.pdf", + "application/pdf", + "content".getBytes(StandardCharsets.UTF_8), + "legacy-request-id", + null + ); + assertThat(missingHeader.statusCode()).isEqualTo(400); + assertThat(JsonPath.read(missingHeader.body(), "$.code")).isEqualTo("VALIDATION_FAILED"); + + HttpResponse invalidHeader = uploadFileWithIdempotencyKey( + rawToken, + "passport.pdf", + "application/pdf", + "content".getBytes(StandardCharsets.UTF_8), + null, + "short" + ); + assertThat(invalidHeader.statusCode()).isEqualTo(400); + assertThat(JsonPath.read(invalidHeader.body(), "$.code")).isEqualTo("VALIDATION_FAILED"); + + HttpResponse invalidDocumentType = uploadFileWithIdempotencyKey( + rawToken, + "passport.pdf", + "application/pdf", + "content".getBytes(StandardCharsets.UTF_8), + null, + "canonical-invalid-document-type", + "CUSTOM_DOCUMENT" + ); + assertThat(invalidDocumentType.statusCode()).isEqualTo(400); + assertThat(JsonPath.read(invalidDocumentType.body(), "$.code")) + .isEqualTo("VALIDATION_FAILED"); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM stored_file WHERE company_id = ?", + Integer.class, + COMPANY_A + )).isZero(); + + HttpResponse first = uploadFileWithIdempotencyKey( + rawToken, + "passport.pdf", + "application/pdf", + "content".getBytes(StandardCharsets.UTF_8), + "shared-legacy-request-id", + "canonical-upload-key-1" + ); + assertThat(first.statusCode()).as(first.body()).isEqualTo(201); + String firstUploadId = JsonPath.read(first.body(), "$.upload_id"); + + HttpResponse retryWithoutLegacyField = uploadFileWithIdempotencyKey( + rawToken, + "passport.pdf", + "application/pdf", + "content".getBytes(StandardCharsets.UTF_8), + null, + "canonical-upload-key-1" + ); + assertThat(retryWithoutLegacyField.statusCode()).as(retryWithoutLegacyField.body()).isEqualTo(201); + assertThat(JsonPath.read(retryWithoutLegacyField.body(), "$.upload_id")) + .isEqualTo(firstUploadId); + + HttpResponse conflict = uploadFileWithIdempotencyKey( + rawToken, + "passport.pdf", + "application/pdf", + "different-content".getBytes(StandardCharsets.UTF_8), + "shared-legacy-request-id", + "canonical-upload-key-1" + ); + assertThat(conflict.statusCode()).isEqualTo(409); + assertThat(JsonPath.read(conflict.body(), "$.code")).isEqualTo("IDEMPOTENCY_CONFLICT"); + + HttpResponse differentCanonicalKey = uploadFileWithIdempotencyKey( + rawToken, + "passport.pdf", + "application/pdf", + "content".getBytes(StandardCharsets.UTF_8), + "shared-legacy-request-id", + "canonical-upload-key-2" + ); + assertThat(differentCanonicalKey.statusCode()).as(differentCanonicalKey.body()).isEqualTo(201); + assertThat(JsonPath.read(differentCanonicalKey.body(), "$.upload_id")) + .isNotEqualTo(firstUploadId); + + UUID workerLinkId = jdbcTemplate.queryForObject( + "SELECT worker_link_id FROM worker_link WHERE task_id = ?", + UUID.class, + UUID.fromString(taskId) + ); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM worker_document_upload_idempotency WHERE worker_link_id = ?", + Integer.class, + workerLinkId + )).isEqualTo(2); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(DISTINCT client_request_id) FROM worker_document_upload_idempotency " + + "WHERE worker_link_id = ? " + + "AND client_request_id = CONCAT('canonical:', CAST(stored_file_id AS VARCHAR)) " + + "AND idempotency_key_hash IS NOT NULL AND request_hash IS NOT NULL", + Integer.class, + workerLinkId + )).isEqualTo(2); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM stored_file WHERE company_id = ?", + Integer.class, + COMPANY_A + )).isEqualTo(2); + } + + @Test + void canonicalUploadDoesNotCollideWithLegacyClientRequestIdEqualToItsKeyHash() throws Exception { + String hrToken = accessToken(login(HR_A_EMAIL)); + String workerId = registerWorker(hrToken, "legacy충돌테스트근로자"); + String taskId = createApprovedTask(hrToken, workerId); + saveDocumentRequestDraft(hrToken, taskId); + String rawToken = issueWorkerLink(hrToken, taskId, "legacy-collision-link-key"); + UUID workerLinkId = jdbcTemplate.queryForObject( + "SELECT worker_link_id FROM worker_link WHERE task_id = ?", + UUID.class, + UUID.fromString(taskId) + ); + String idempotencyKey = "legacy-collision-upload-key"; + String idempotencyKeyHash = workerLinkHasher.hash(idempotencyKey); + UUID legacyStoredFileId = UUID.randomUUID(); + + jdbcTemplate.update( + """ + INSERT INTO stored_file ( + stored_file_id, company_id, name, mime_type, size, purpose, + task_id, storage_key, scan_status, verified + ) VALUES (?, ?, 'legacy.pdf', 'application/pdf', 1, 'WORKER_LINK_SUBMISSION', + ?, ?, 'NOT_SCANNED', FALSE) + """, + legacyStoredFileId, + COMPANY_A, + UUID.fromString(taskId), + "legacy-collision-" + legacyStoredFileId + ); + jdbcTemplate.update( + """ + INSERT INTO worker_document_upload_idempotency ( + worker_link_id, company_id, client_request_id, stored_file_id + ) VALUES (?, ?, ?, ?) + """, + workerLinkId, + COMPANY_A, + idempotencyKeyHash, + legacyStoredFileId + ); + + HttpResponse response = uploadFileWithIdempotencyKey( + rawToken, + "passport.pdf", + "application/pdf", + "canonical-content".getBytes(StandardCharsets.UTF_8), + null, + idempotencyKey + ); + + assertThat(response.statusCode()).as(response.body()).isEqualTo(201); + UUID canonicalStoredFileId = UUID.fromString(JsonPath.read(response.body(), "$.upload_id")); + assertThat(jdbcTemplate.queryForObject( + "SELECT client_request_id FROM worker_document_upload_idempotency " + + "WHERE worker_link_id = ? AND idempotency_key_hash = ?", + String.class, + workerLinkId, + idempotencyKeyHash + )).isEqualTo("canonical:" + canonicalStoredFileId); + assertThat(jdbcTemplate.queryForObject( + "SELECT COUNT(*) FROM worker_document_upload_idempotency WHERE worker_link_id = ?", + Integer.class, + workerLinkId + )).isEqualTo(2); + } + @Test void requestedSlotAnswerIsStoredAndDifferentIdempotentPayloadIsRejected() throws Exception { String hrToken = accessToken(login(HR_A_EMAIL)); @@ -1078,6 +1271,28 @@ void documentsWorkerLinkAndWorkerActivityEndpointsInOpenApi() throws Exception { response.body(), "$.paths['/api/v1/workers/{workerId}/activities'].get.operationId" )).isEqualTo("getWorkerActivities"); + assertThat(JsonPath.read( + response.body(), + "$.paths['/api/v1/public/worker-links/{token}/documents'].post.operationId" + )).isEqualTo("uploadWorkerLinkDocument"); + assertThat(JsonPath.>read( + response.body(), + "$.paths['/api/v1/public/worker-links/{token}/documents'].post.parameters" + + "[?(@.name == 'Idempotency-Key')].required" + )).containsExactly(true); + assertThat(JsonPath.>>read( + response.body(), + "$.paths['/api/v1/public/worker-links/{token}/documents'].post.parameters" + + "[?(@.name == 'documentType')].schema.enum" + )).containsExactly(List.of( + "PASSPORT_COPY", + "ARC", + "CONTRACT", + "PERMIT", + "EMPLOYMENT_EXTENSION_APPLICATION", + "INTEGRATED_APPLICATION", + "RESIDENCE_PROOF" + )); assertThat(JsonPath.read( response.body(), "$.components.schemas.WorkerLinkIssueRequest.properties.expires_in_hours.minimum" @@ -1264,6 +1479,7 @@ private HttpResponse uploadFile(String token, String filename, String mi HttpRequest request = HttpRequest.newBuilder(uri("/api/v1/public/worker-links/" + token + "/documents")) .header(HttpHeaders.CONTENT_TYPE, "multipart/form-data; boundary=" + BOUNDARY) + .header("Idempotency-Key", "upload-" + UUID.randomUUID()) .POST(HttpRequest.BodyPublishers.ofByteArray(out.toByteArray())) .build(); return httpClient.send(request, HttpResponse.BodyHandlers.ofString()); @@ -1284,24 +1500,63 @@ private HttpResponse uploadFileAsType( HttpRequest request = HttpRequest.newBuilder(uri("/api/v1/public/worker-links/" + token + "/documents")) .header(HttpHeaders.CONTENT_TYPE, "multipart/form-data; boundary=" + BOUNDARY) + .header("Idempotency-Key", "upload-" + UUID.randomUUID()) .POST(HttpRequest.BodyPublishers.ofByteArray(out.toByteArray())) .build(); return httpClient.send(request, HttpResponse.BodyHandlers.ofString()); } - private HttpResponse uploadFileWithFixedClientRequestId( - String token, String filename, String mimeType, byte[] content, String clientRequestId + private HttpResponse uploadFileWithIdempotencyKey( + String token, + String filename, + String mimeType, + byte[] content, + String clientRequestId, + String idempotencyKey + ) throws Exception { + return uploadFileWithIdempotencyKey( + token, + filename, + mimeType, + content, + clientRequestId, + idempotencyKey, + null + ); + } + + private HttpResponse uploadFileWithIdempotencyKey( + String token, + String filename, + String mimeType, + byte[] content, + String clientRequestId, + String idempotencyKey, + String documentType ) throws Exception { ByteArrayOutputStream out = new ByteArrayOutputStream(); writePart(out, "file", filename, mimeType, content); - writeFieldPart(out, "clientRequestId", clientRequestId); + if (clientRequestId != null) { + writeFieldPart(out, "clientRequestId", clientRequestId); + } + if (documentType != null) { + writeFieldPart(out, "documentType", documentType); + } out.write(("--" + BOUNDARY + "--\r\n").getBytes(StandardCharsets.UTF_8)); - HttpRequest request = HttpRequest.newBuilder(uri("/api/v1/public/worker-links/" + token + "/documents")) - .header(HttpHeaders.CONTENT_TYPE, "multipart/form-data; boundary=" + BOUNDARY) + HttpRequest.Builder request = HttpRequest.newBuilder( + uri("/api/v1/public/worker-links/" + token + "/documents") + ) + .header(HttpHeaders.CONTENT_TYPE, "multipart/form-data; boundary=" + BOUNDARY); + if (idempotencyKey != null) { + request.header("Idempotency-Key", idempotencyKey); + } + return httpClient.send( + request .POST(HttpRequest.BodyPublishers.ofByteArray(out.toByteArray())) - .build(); - return httpClient.send(request, HttpResponse.BodyHandlers.ofString()); + .build(), + HttpResponse.BodyHandlers.ofString() + ); } private void writePart(ByteArrayOutputStream out, String name, String filename, String mimeType, byte[] content)