diff --git a/src/main/java/com/saferoute/domain/congestion/controller/CongestionController.java b/src/main/java/com/saferoute/domain/congestion/controller/CongestionController.java index 7230add3..dcbacb8e 100644 --- a/src/main/java/com/saferoute/domain/congestion/controller/CongestionController.java +++ b/src/main/java/com/saferoute/domain/congestion/controller/CongestionController.java @@ -1,8 +1,10 @@ package com.saferoute.domain.congestion.controller; import com.saferoute.domain.congestion.dto.request.ReportCongestionRequest; +import com.saferoute.domain.congestion.dto.request.ConnectEventImageRequest; import com.saferoute.domain.congestion.dto.response.ObservationResponse; import com.saferoute.domain.congestion.service.CongestionEventService; +import com.saferoute.domain.congestion.service.CongestionEventImageService; import com.saferoute.domain.device.service.DeviceAuthorizationService; import com.saferoute.domain.telemetry.dynamo.entity.ObservationItem; import com.saferoute.domain.telemetry.dynamo.repository.IdempotentSaveResult; @@ -14,9 +16,12 @@ import org.springframework.http.ResponseEntity; import org.springframework.security.core.annotation.AuthenticationPrincipal; import org.springframework.web.bind.annotation.PostMapping; +import org.springframework.web.bind.annotation.PatchMapping; +import org.springframework.web.bind.annotation.PathVariable; import org.springframework.web.bind.annotation.RequestBody; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.RestController; +import java.util.UUID; @Tag(name = "혼잡도", description = "CCTV 혼잡 이벤트 수신 API") @RestController @@ -25,6 +30,7 @@ public class CongestionController { private final CongestionEventService congestionEventService; + private final CongestionEventImageService congestionEventImageService; private final DeviceAuthorizationService deviceAuthorizationService; @PostMapping @@ -38,4 +44,14 @@ public ResponseEntity reportCongestion( HttpStatus status = saveResult.created() ? HttpStatus.CREATED : HttpStatus.OK; return ResponseEntity.status(status).body(response); } + + @PatchMapping("/{eventId}/image") + public ResponseEntity connectEventImage( + @AuthenticationPrincipal DevicePrincipal principal, + @PathVariable UUID eventId, + @Valid @RequestBody ConnectEventImageRequest request + ) { + congestionEventImageService.connectImage(principal, eventId, request); + return ResponseEntity.noContent().build(); + } } diff --git a/src/main/java/com/saferoute/domain/congestion/dto/request/ConnectEventImageRequest.java b/src/main/java/com/saferoute/domain/congestion/dto/request/ConnectEventImageRequest.java new file mode 100644 index 00000000..e5185b6a --- /dev/null +++ b/src/main/java/com/saferoute/domain/congestion/dto/request/ConnectEventImageRequest.java @@ -0,0 +1,11 @@ +package com.saferoute.domain.congestion.dto.request; + +import jakarta.validation.constraints.NotBlank; +import jakarta.validation.constraints.NotNull; +import jakarta.validation.constraints.PositiveOrZero; + +public record ConnectEventImageRequest( + @NotBlank String eventImageKey, + @NotNull @PositiveOrZero Long uploadedAt +) { +} diff --git a/src/main/java/com/saferoute/domain/congestion/service/CongestionEventImageService.java b/src/main/java/com/saferoute/domain/congestion/service/CongestionEventImageService.java new file mode 100644 index 00000000..01ce9adf --- /dev/null +++ b/src/main/java/com/saferoute/domain/congestion/service/CongestionEventImageService.java @@ -0,0 +1,142 @@ +package com.saferoute.domain.congestion.service; + +import com.saferoute.domain.congestion.dto.request.ConnectEventImageRequest; +import com.saferoute.domain.device.service.DeviceAuthorizationService; +import com.saferoute.domain.telemetry.dynamo.entity.EventProcessingStatus; +import com.saferoute.domain.telemetry.dynamo.entity.ImageUploadStatus; +import com.saferoute.domain.telemetry.dynamo.entity.ObservationItem; +import com.saferoute.domain.telemetry.dynamo.repository.ObservationRepository; +import com.saferoute.global.api.error.CongestionErrorCode; +import com.saferoute.global.api.exception.ApiException; +import com.saferoute.global.security.DevicePrincipal; +import com.saferoute.infrastructure.s3.service.S3Service; +import com.saferoute.infrastructure.websocket.service.TrainingEventPublisher; +import java.util.Objects; +import java.util.UUID; +import lombok.RequiredArgsConstructor; +import org.springframework.stereotype.Service; +import org.springframework.transaction.annotation.Transactional; + +@Service +@RequiredArgsConstructor +@Transactional(readOnly = true) +public class CongestionEventImageService { + + private static final String TRAINING_PREFIX = "training"; + private static final String EVENTS_DIRECTORY = "events"; + private static final String JPEG_SUFFIX = ".jpg"; + + private final ObservationRepository observationRepository; + private final DeviceAuthorizationService deviceAuthorizationService; + private final S3Service s3Service; + private final TrainingEventPublisher trainingEventPublisher; + + public void connectImage( + DevicePrincipal principal, + UUID eventId, + ConnectEventImageRequest request + ) { + ObservationItem item = findEvent(eventId); + deviceAuthorizationService.validateCctv(principal, item.getCctvCode()); + validateEventProcessed(item); + validateObjectKey(item, eventId, request.eventImageKey()); + + if (isSameCompletedImage(item, request)) { + publishImageUpdated(item); + return; + } + validateImageState(item); + + if (!s3Service.objectExists(request.eventImageKey())) { + throw new ApiException(CongestionErrorCode.EVENT_IMAGE_OBJECT_NOT_FOUND); + } + + if (!observationRepository.completeImageUpload( + eventId.toString(), request.eventImageKey(), request.uploadedAt() + )) { + ObservationItem latest = findEvent(eventId); + if (!isSameCompletedImage(latest, request)) { + throw new ApiException(CongestionErrorCode.EVENT_IMAGE_STATE_CONFLICT); + } + publishImageUpdated(latest); + return; + } + + item.setEventImageKey(request.eventImageKey()); + item.setImageUploadedAt(request.uploadedAt()); + item.setImageUploadStatus(ImageUploadStatus.COMPLETED); + publishImageUpdated(item); + } + + private ObservationItem findEvent(UUID eventId) { + return observationRepository.findByEventId(eventId.toString()) + .orElseThrow(() -> new ApiException(CongestionErrorCode.EVENT_NOT_FOUND)); + } + + private void validateEventProcessed(ObservationItem item) { + if (item.getEventStatus() != EventProcessingStatus.PROCESSED) { + throw new ApiException(CongestionErrorCode.EVENT_NOT_PROCESSED); + } + } + + private void validateImageState(ObservationItem item) { + if (item.getImageUploadStatus() != null + && item.getImageUploadStatus() != ImageUploadStatus.PENDING + && item.getImageUploadStatus() != ImageUploadStatus.FAILED) { + throw new ApiException(CongestionErrorCode.EVENT_IMAGE_STATE_CONFLICT); + } + } + + private void validateObjectKey(ObservationItem item, UUID eventId, String objectKey) { + String[] segments = objectKey.split("/", -1); + if (segments.length != 5 + || !TRAINING_PREFIX.equals(segments[0]) + || !EVENTS_DIRECTORY.equals(segments[2]) + || !segments[4].endsWith(JPEG_SUFFIX)) { + throw new ApiException(CongestionErrorCode.EVENT_IMAGE_KEY_INVALID); + } + + UUID sessionId = parseCanonicalUuid(segments[1]); + String imageEventId = segments[4].substring(0, segments[4].length() - JPEG_SUFFIX.length()); + UUID keyEventId = parseCanonicalUuid(imageEventId); + boolean sameIdentity = Objects.equals(item.getTrainingSessionId(), sessionId.toString()) + && Objects.equals(item.getCctvCode(), segments[3]) + && eventId.equals(keyEventId); + if (!sameIdentity) { + throw new ApiException(CongestionErrorCode.EVENT_IMAGE_IDENTITY_MISMATCH); + } + } + + private UUID parseCanonicalUuid(String value) { + try { + UUID uuid = UUID.fromString(value); + if (!uuid.toString().equals(value)) { + throw new IllegalArgumentException("non-canonical UUID"); + } + return uuid; + } catch (IllegalArgumentException exception) { + throw new ApiException(CongestionErrorCode.EVENT_IMAGE_KEY_INVALID, exception); + } + } + + private boolean isSameCompletedImage( + ObservationItem item, + ConnectEventImageRequest request + ) { + return isCompletedStatus(item.getImageUploadStatus()) + && Objects.equals(item.getEventImageKey(), request.eventImageKey()) + && Objects.equals(item.getImageUploadedAt(), request.uploadedAt()); + } + + @SuppressWarnings("deprecation") + private boolean isCompletedStatus(ImageUploadStatus status) { + return status == ImageUploadStatus.COMPLETED || status == ImageUploadStatus.UPLOADED; + } + + private void publishImageUpdated(ObservationItem item) { + trainingEventPublisher.publishCongestionImageUpdated( + UUID.fromString(item.getTrainingSessionId()), + item + ); + } +} diff --git a/src/main/java/com/saferoute/domain/telemetry/dynamo/entity/ImageUploadStatus.java b/src/main/java/com/saferoute/domain/telemetry/dynamo/entity/ImageUploadStatus.java index 1e65a565..c32aa3ed 100644 --- a/src/main/java/com/saferoute/domain/telemetry/dynamo/entity/ImageUploadStatus.java +++ b/src/main/java/com/saferoute/domain/telemetry/dynamo/entity/ImageUploadStatus.java @@ -2,6 +2,10 @@ public enum ImageUploadStatus { PENDING, - UPLOADED, - FAILED + COMPLETED, + FAILED, + + /** 기존 DynamoDB 항목 역직렬화를 위한 호환 값. 신규 저장에는 사용하지 않는다. */ + @Deprecated + UPLOADED } diff --git a/src/main/java/com/saferoute/domain/telemetry/dynamo/entity/ObservationItem.java b/src/main/java/com/saferoute/domain/telemetry/dynamo/entity/ObservationItem.java index 20975a9e..410f06bb 100644 --- a/src/main/java/com/saferoute/domain/telemetry/dynamo/entity/ObservationItem.java +++ b/src/main/java/com/saferoute/domain/telemetry/dynamo/entity/ObservationItem.java @@ -32,6 +32,9 @@ public class ObservationItem { private Long windowEnd; private Long capturedAt; private String monitoringImageKey; + private String eventImageKey; + private Long imageUploadedAt; + private ImageUploadStatus imageUploadStatus; private Long configVersion; private Long expiresAt; private EventProcessingStatus eventStatus; @@ -72,6 +75,7 @@ public static ObservationItem create( item.windowEnd = windowEnd; item.capturedAt = capturedAt; item.monitoringImageKey = monitoringImageKey; + item.imageUploadStatus = ImageUploadStatus.PENDING; item.configVersion = configVersion; item.expiresAt = Math.floorDiv(capturedAt, 1_000L) + TTL_SECONDS; item.eventStatus = EventProcessingStatus.RECEIVED; @@ -236,6 +240,30 @@ public void setMonitoringImageKey(String monitoringImageKey) { this.monitoringImageKey = monitoringImageKey; } + public String getEventImageKey() { + return eventImageKey; + } + + public void setEventImageKey(String eventImageKey) { + this.eventImageKey = eventImageKey; + } + + public Long getImageUploadedAt() { + return imageUploadedAt; + } + + public void setImageUploadedAt(Long imageUploadedAt) { + this.imageUploadedAt = imageUploadedAt; + } + + public ImageUploadStatus getImageUploadStatus() { + return imageUploadStatus; + } + + public void setImageUploadStatus(ImageUploadStatus imageUploadStatus) { + this.imageUploadStatus = imageUploadStatus; + } + public Long getConfigVersion() { return configVersion; } diff --git a/src/main/java/com/saferoute/domain/telemetry/dynamo/repository/CongestionEventRepository.java b/src/main/java/com/saferoute/domain/telemetry/dynamo/repository/CongestionEventRepository.java index 790d993e..d6a0de98 100644 --- a/src/main/java/com/saferoute/domain/telemetry/dynamo/repository/CongestionEventRepository.java +++ b/src/main/java/com/saferoute/domain/telemetry/dynamo/repository/CongestionEventRepository.java @@ -162,12 +162,15 @@ private void validateEventStatusTransition( } } + @SuppressWarnings("deprecation") private void validateImageStatusTransition( ImageUploadStatus expectedStatus, ImageUploadStatus newStatus ) { boolean complete = expectedStatus == ImageUploadStatus.PENDING - && (newStatus == ImageUploadStatus.UPLOADED || newStatus == ImageUploadStatus.FAILED); + && (newStatus == ImageUploadStatus.COMPLETED + || newStatus == ImageUploadStatus.UPLOADED + || newStatus == ImageUploadStatus.FAILED); boolean retry = expectedStatus == ImageUploadStatus.FAILED && newStatus == ImageUploadStatus.PENDING; if (!complete && !retry) { diff --git a/src/main/java/com/saferoute/domain/telemetry/dynamo/repository/ObservationRepository.java b/src/main/java/com/saferoute/domain/telemetry/dynamo/repository/ObservationRepository.java index 2624e1e1..7e5e73e9 100644 --- a/src/main/java/com/saferoute/domain/telemetry/dynamo/repository/ObservationRepository.java +++ b/src/main/java/com/saferoute/domain/telemetry/dynamo/repository/ObservationRepository.java @@ -1,6 +1,7 @@ package com.saferoute.domain.telemetry.dynamo.repository; import com.saferoute.domain.telemetry.dynamo.entity.EventProcessingStatus; +import com.saferoute.domain.telemetry.dynamo.entity.ImageUploadStatus; import com.saferoute.domain.telemetry.dynamo.entity.ObservationItem; import java.util.List; import java.util.Optional; @@ -114,6 +115,26 @@ public boolean failProcessing(String eventId, String processingOwner) { ); } + public boolean completeImageUpload(String eventId, String eventImageKey, long uploadedAt) { + ObservationItem item = processingUpdateItem(eventId); + item.setEventImageKey(eventImageKey); + item.setImageUploadedAt(uploadedAt); + item.setImageUploadStatus(ImageUploadStatus.COMPLETED); + + Expression condition = Expression.builder() + .expression("attribute_exists(#pk) AND #eventStatus = :processed" + + " AND (attribute_not_exists(#imageStatus)" + + " OR #imageStatus = :pending OR #imageStatus = :failed)") + .putExpressionName("#pk", "pk") + .putExpressionName("#eventStatus", "eventStatus") + .putExpressionName("#imageStatus", "imageUploadStatus") + .putExpressionValue(":processed", AttributeValue.fromS("PROCESSED")) + .putExpressionValue(":pending", AttributeValue.fromS("PENDING")) + .putExpressionValue(":failed", AttributeValue.fromS("FAILED")) + .build(); + return updateConditionally(item, condition); + } + public List findAllBySessionIdAndCctvCode( String trainingSessionId, String cctvCode diff --git a/src/main/java/com/saferoute/global/api/error/CongestionErrorCode.java b/src/main/java/com/saferoute/global/api/error/CongestionErrorCode.java index fda750c2..974b729b 100644 --- a/src/main/java/com/saferoute/global/api/error/CongestionErrorCode.java +++ b/src/main/java/com/saferoute/global/api/error/CongestionErrorCode.java @@ -18,6 +18,36 @@ public enum CongestionErrorCode implements BaseErrorCode { HttpStatus.CONFLICT, "CONGESTION002", "동일한 eventId에 다른 세션, CCTV 또는 경로 정보가 전달되었습니다." + ), + EVENT_NOT_FOUND( + HttpStatus.NOT_FOUND, + "CONGESTION003", + "혼잡 이벤트를 찾을 수 없습니다." + ), + EVENT_NOT_PROCESSED( + HttpStatus.CONFLICT, + "CONGESTION004", + "처리가 완료된 혼잡 이벤트에만 이미지를 연결할 수 있습니다." + ), + EVENT_IMAGE_STATE_CONFLICT( + HttpStatus.CONFLICT, + "CONGESTION005", + "현재 이미지 상태에서는 이미지를 연결할 수 없습니다." + ), + EVENT_IMAGE_KEY_INVALID( + HttpStatus.BAD_REQUEST, + "CONGESTION006", + "이벤트 이미지 경로 형식이 올바르지 않습니다." + ), + EVENT_IMAGE_IDENTITY_MISMATCH( + HttpStatus.CONFLICT, + "CONGESTION007", + "이미지 경로의 세션, CCTV 또는 eventId가 이벤트와 일치하지 않습니다." + ), + EVENT_IMAGE_OBJECT_NOT_FOUND( + HttpStatus.CONFLICT, + "CONGESTION008", + "업로드가 완료된 이벤트 이미지 객체를 찾을 수 없습니다." ); private final HttpStatus httpStatus; diff --git a/src/main/java/com/saferoute/global/api/error/S3ErrorCode.java b/src/main/java/com/saferoute/global/api/error/S3ErrorCode.java index 557d8864..8e45fba6 100644 --- a/src/main/java/com/saferoute/global/api/error/S3ErrorCode.java +++ b/src/main/java/com/saferoute/global/api/error/S3ErrorCode.java @@ -15,6 +15,11 @@ public enum S3ErrorCode implements BaseErrorCode { HttpStatus.INTERNAL_SERVER_ERROR, "S3_ERROR_003", "S3 업로드 URL 발급에 실패했습니다." + ), + OBJECT_CHECK_FAILED( + HttpStatus.SERVICE_UNAVAILABLE, + "S3_ERROR_004", + "S3 객체 확인에 실패했습니다. 잠시 후 다시 시도해 주세요." ); private final HttpStatus httpStatus; diff --git a/src/main/java/com/saferoute/infrastructure/s3/service/S3Service.java b/src/main/java/com/saferoute/infrastructure/s3/service/S3Service.java index 52c0c6d8..e098bcd5 100644 --- a/src/main/java/com/saferoute/infrastructure/s3/service/S3Service.java +++ b/src/main/java/com/saferoute/infrastructure/s3/service/S3Service.java @@ -13,6 +13,9 @@ import software.amazon.awssdk.core.sync.RequestBody; import software.amazon.awssdk.services.s3.S3Client; import software.amazon.awssdk.services.s3.model.PutObjectRequest; +import software.amazon.awssdk.services.s3.model.HeadObjectRequest; +import software.amazon.awssdk.services.s3.model.NoSuchKeyException; +import software.amazon.awssdk.services.s3.model.S3Exception; @Service public class S3Service { @@ -59,6 +62,26 @@ public S3UploadResponse upload(MultipartFile file) { ); } + public boolean objectExists(String key) { + HeadObjectRequest request = HeadObjectRequest.builder() + .bucket(s3Properties.bucket()) + .key(key) + .build(); + try { + s3Client.headObject(request); + return true; + } catch (NoSuchKeyException exception) { + return false; + } catch (S3Exception exception) { + if (exception.statusCode() == 404) { + return false; + } + throw new ApiException(S3ErrorCode.OBJECT_CHECK_FAILED, exception); + } catch (SdkException exception) { + throw new ApiException(S3ErrorCode.OBJECT_CHECK_FAILED, exception); + } + } + private String createObjectKey(String originalFilename) { String filename = originalFilename == null ? "file" : originalFilename; filename = filename.replace('\\', '/'); diff --git a/src/main/java/com/saferoute/infrastructure/websocket/dto/CongestionImageUpdatedData.java b/src/main/java/com/saferoute/infrastructure/websocket/dto/CongestionImageUpdatedData.java new file mode 100644 index 00000000..c920e117 --- /dev/null +++ b/src/main/java/com/saferoute/infrastructure/websocket/dto/CongestionImageUpdatedData.java @@ -0,0 +1,21 @@ +package com.saferoute.infrastructure.websocket.dto; + +import com.saferoute.domain.telemetry.dynamo.entity.ImageUploadStatus; +import com.saferoute.domain.telemetry.dynamo.entity.ObservationItem; +import java.util.UUID; + +public record CongestionImageUpdatedData( + UUID eventId, + String eventImageKey, + Long uploadedAt, + ImageUploadStatus imageUploadStatus +) { + public static CongestionImageUpdatedData from(ObservationItem item) { + return new CongestionImageUpdatedData( + UUID.fromString(item.getEventId()), + item.getEventImageKey(), + item.getImageUploadedAt(), + item.getImageUploadStatus() + ); + } +} diff --git a/src/main/java/com/saferoute/infrastructure/websocket/dto/TrainingEventType.java b/src/main/java/com/saferoute/infrastructure/websocket/dto/TrainingEventType.java index 1e2d6b5c..e023caa0 100644 --- a/src/main/java/com/saferoute/infrastructure/websocket/dto/TrainingEventType.java +++ b/src/main/java/com/saferoute/infrastructure/websocket/dto/TrainingEventType.java @@ -9,6 +9,9 @@ public enum TrainingEventType { // CongestionEventService.reportCongestion() 호출 시 발행된다. (이슈 #48) CONGESTION_UPDATED, + // 혼잡 이벤트 이미지가 S3 업로드 완료 후 연결되었을 때 발행된다. (이슈 #83) + CONGESTION_IMAGE_UPDATED, + // RouteRecalculationService.trigger()가 재탐색 결과를 PENDING으로 저장했을 때 발행된다. (이슈 #48) ROUTE_RECALCULATION_REQUESTED, diff --git a/src/main/java/com/saferoute/infrastructure/websocket/service/TrainingEventPublisher.java b/src/main/java/com/saferoute/infrastructure/websocket/service/TrainingEventPublisher.java index b60baead..908dd632 100644 --- a/src/main/java/com/saferoute/infrastructure/websocket/service/TrainingEventPublisher.java +++ b/src/main/java/com/saferoute/infrastructure/websocket/service/TrainingEventPublisher.java @@ -6,6 +6,7 @@ import com.saferoute.domain.telemetry.dynamo.entity.ObservationItem; import com.saferoute.domain.training.entity.TrainingSession; import com.saferoute.infrastructure.websocket.dto.CongestionEventData; +import com.saferoute.infrastructure.websocket.dto.CongestionImageUpdatedData; import com.saferoute.infrastructure.websocket.dto.IoTLightEventMessage; import com.saferoute.infrastructure.websocket.dto.IoTLightStatusEventData; import com.saferoute.infrastructure.websocket.dto.RouteRecalculationEventData; @@ -149,6 +150,21 @@ public void publishCongestionUpdated(UUID sessionId, UUID edgeId, ObservationIte ); } + public void publishCongestionImageUpdated(UUID sessionId, ObservationItem item) { + TrainingEventMessage message = TrainingEventMessage.of( + TrainingEventType.CONGESTION_IMAGE_UPDATED, + sessionId, + CongestionImageUpdatedData.from(item) + ); + + messagingTemplate.convertAndSend(SESSION_TOPIC_PREFIX + sessionId, message); + log.debug( + "혼잡 이벤트 이미지 갱신 발행: sessionId={}, eventId={}", + sessionId, + item.getEventId() + ); + } + // DB 트랜잭션이 실제로 커밋된 이후에만 발행한다 (publishTrainingStatusUpdatedAfterCommit 참고). public void publishRouteRecalculationRequestedAfterCommit(RouteRecalculation recalculation) { publishAfterCommit( diff --git a/src/test/java/com/saferoute/domain/congestion/controller/CongestionControllerTest.java b/src/test/java/com/saferoute/domain/congestion/controller/CongestionControllerTest.java index 02497f7f..d2b9ba89 100644 --- a/src/test/java/com/saferoute/domain/congestion/controller/CongestionControllerTest.java +++ b/src/test/java/com/saferoute/domain/congestion/controller/CongestionControllerTest.java @@ -2,7 +2,10 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.BDDMockito.given; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.doThrow; import static org.springframework.test.web.servlet.request.MockMvcRequestBuilders.post; +import static org.springframework.test.web.servlet.request.MockMvcRequestBuilders.patch; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.jsonPath; import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; @@ -11,6 +14,7 @@ import com.saferoute.domain.congestion.dto.request.ReportCongestionRequest; import com.saferoute.domain.congestion.entity.CongestionLevel; import com.saferoute.domain.congestion.service.CongestionEventService; +import com.saferoute.domain.congestion.service.CongestionEventImageService; import com.saferoute.domain.device.service.DeviceAuthorizationService; import com.saferoute.domain.telemetry.dynamo.entity.ObservationItem; import com.saferoute.domain.telemetry.dynamo.repository.IdempotentSaveResult; @@ -60,6 +64,9 @@ class CongestionControllerTest { @MockitoBean private CongestionEventService congestionEventService; + @MockitoBean + private CongestionEventImageService congestionEventImageService; + @MockitoBean private DeviceAuthorizationService deviceAuthorizationService; @@ -253,6 +260,69 @@ void reportCongestion_returnsConflictWhenEventIdentityMismatches() throws Except .andExpect(jsonPath("$.code").value("CONGESTION002")); } + @Test + @DisplayName("이벤트 이미지 연결 성공 시 204를 반환한다") + void connectEventImage_returnsNoContent() throws Exception { + ObjectNode body = objectMapper.createObjectNode() + .put("eventImageKey", "training/" + SESSION_ID + "/events/CCTV_001/" + + OBSERVATION_ID + ".jpg") + .put("uploadedAt", 1_786_500_002_800L); + + mockMvc.perform(patch("/api/v1/device/congestion-events/{eventId}/image", OBSERVATION_ID) + .contentType(MediaType.APPLICATION_JSON) + .content(objectMapper.writeValueAsString(body))) + .andExpect(status().isNoContent()); + + verify(congestionEventImageService).connectImage(any(), any(), any()); + } + + @Test + @DisplayName("eventImageKey가 없으면 400을 반환한다") + void connectEventImage_returnsBadRequestWithoutImageKey() throws Exception { + ObjectNode body = objectMapper.createObjectNode() + .put("uploadedAt", 1_786_500_002_800L); + + mockMvc.perform(patch("/api/v1/device/congestion-events/{eventId}/image", OBSERVATION_ID) + .contentType(MediaType.APPLICATION_JSON) + .content(objectMapper.writeValueAsString(body))) + .andExpect(status().isBadRequest()); + } + + @Test + @DisplayName("이벤트를 찾을 수 없으면 CONGESTION003과 404를 반환한다") + void connectEventImage_returnsNotFoundWhenEventMissing() throws Exception { + doThrow(new ApiException(CongestionErrorCode.EVENT_NOT_FOUND)) + .when(congestionEventImageService).connectImage(any(), any(), any()); + + mockMvc.perform(patch("/api/v1/device/congestion-events/{eventId}/image", OBSERVATION_ID) + .contentType(MediaType.APPLICATION_JSON) + .content(objectMapper.writeValueAsString(validImageRequest()))) + .andExpect(status().isNotFound()) + .andExpect(jsonPath("$.isSuccess").value(false)) + .andExpect(jsonPath("$.code").value("CONGESTION003")); + } + + @Test + @DisplayName("이미지 상태가 충돌하면 CONGESTION005와 409를 반환한다") + void connectEventImage_returnsConflictWhenImageStateConflicts() throws Exception { + doThrow(new ApiException(CongestionErrorCode.EVENT_IMAGE_STATE_CONFLICT)) + .when(congestionEventImageService).connectImage(any(), any(), any()); + + mockMvc.perform(patch("/api/v1/device/congestion-events/{eventId}/image", OBSERVATION_ID) + .contentType(MediaType.APPLICATION_JSON) + .content(objectMapper.writeValueAsString(validImageRequest()))) + .andExpect(status().isConflict()) + .andExpect(jsonPath("$.isSuccess").value(false)) + .andExpect(jsonPath("$.code").value("CONGESTION005")); + } + + private ObjectNode validImageRequest() { + return objectMapper.createObjectNode() + .put("eventImageKey", "training/" + SESSION_ID + "/events/CCTV_001/" + + OBSERVATION_ID + ".jpg") + .put("uploadedAt", 1_786_500_002_800L); + } + private ObservationItem observation() { return ObservationItem.create( OBSERVATION_ID, SESSION_ID, EDGE_ID, "CCTV_001", 5.0, 8, 25, diff --git a/src/test/java/com/saferoute/domain/congestion/service/CongestionEventImageServiceTest.java b/src/test/java/com/saferoute/domain/congestion/service/CongestionEventImageServiceTest.java new file mode 100644 index 00000000..89dc5396 --- /dev/null +++ b/src/test/java/com/saferoute/domain/congestion/service/CongestionEventImageServiceTest.java @@ -0,0 +1,247 @@ +package com.saferoute.domain.congestion.service; + +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.BDDMockito.given; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; + +import com.saferoute.domain.congestion.dto.request.ConnectEventImageRequest; +import com.saferoute.domain.congestion.entity.CongestionLevel; +import com.saferoute.domain.device.service.DeviceAuthorizationService; +import com.saferoute.domain.telemetry.dynamo.entity.EventProcessingStatus; +import com.saferoute.domain.telemetry.dynamo.entity.ImageUploadStatus; +import com.saferoute.domain.telemetry.dynamo.entity.ObservationItem; +import com.saferoute.domain.telemetry.dynamo.repository.ObservationRepository; +import com.saferoute.global.api.error.CongestionErrorCode; +import com.saferoute.global.api.error.DeviceErrorCode; +import com.saferoute.global.api.exception.ApiException; +import com.saferoute.global.security.DevicePrincipal; +import com.saferoute.infrastructure.s3.service.S3Service; +import com.saferoute.infrastructure.websocket.service.TrainingEventPublisher; +import java.util.Optional; +import java.util.UUID; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +@ExtendWith(MockitoExtension.class) +class CongestionEventImageServiceTest { + + private static final UUID EVENT_ID = UUID.fromString("00000000-0000-0000-0000-000000000001"); + private static final UUID SESSION_ID = UUID.fromString("00000000-0000-0000-0000-000000000002"); + private static final UUID EDGE_ID = UUID.fromString("00000000-0000-0000-0000-000000000003"); + private static final String CCTV_CODE = "CCTV_001"; + private static final String IMAGE_KEY = "training/" + SESSION_ID + "/events/" + + CCTV_CODE + "/" + EVENT_ID + ".jpg"; + + @InjectMocks + private CongestionEventImageService service; + @Mock + private ObservationRepository observationRepository; + @Mock + private DeviceAuthorizationService deviceAuthorizationService; + @Mock + private S3Service s3Service; + @Mock + private TrainingEventPublisher trainingEventPublisher; + + private final DevicePrincipal principal = new DevicePrincipal(UUID.randomUUID(), CCTV_CODE); + private final ConnectEventImageRequest request = new ConnectEventImageRequest(IMAGE_KEY, 1_786_500_002_800L); + private ObservationItem item; + + @BeforeEach + void setUp() { + item = ObservationItem.create( + EVENT_ID, SESSION_ID, EDGE_ID, CCTV_CODE, 5.0, 8, 25, 2.5, + CongestionLevel.CROWDED, 1_000L, 2_000L, 2_000L, null, 1L + ); + item.setEventStatus(EventProcessingStatus.PROCESSED); + } + + @Test + void 처리된_이벤트에_S3_이미지를_연결하고_업데이트를_발행한다() { + given(observationRepository.findByEventId(EVENT_ID.toString())).willReturn(Optional.of(item)); + given(s3Service.objectExists(IMAGE_KEY)).willReturn(true); + given(observationRepository.completeImageUpload( + EVENT_ID.toString(), IMAGE_KEY, request.uploadedAt() + )).willReturn(true); + + service.connectImage(principal, EVENT_ID, request); + + verify(deviceAuthorizationService).validateCctv(principal, CCTV_CODE); + verify(observationRepository).completeImageUpload( + EVENT_ID.toString(), IMAGE_KEY, request.uploadedAt() + ); + verify(trainingEventPublisher).publishCongestionImageUpdated(SESSION_ID, item); + } + + @Test + void 이벤트_POST보다_PATCH가_먼저_도착하면_404를_반환한다() { + given(observationRepository.findByEventId(EVENT_ID.toString())).willReturn(Optional.empty()); + + assertError(CongestionErrorCode.EVENT_NOT_FOUND, request); + + verify(s3Service, never()).objectExists(IMAGE_KEY); + } + + @Test + void 이벤트가_PROCESSED가_아니면_이미지를_연결하지_않는다() { + item.setEventStatus(EventProcessingStatus.PROCESSING); + given(observationRepository.findByEventId(EVENT_ID.toString())).willReturn(Optional.of(item)); + + assertError(CongestionErrorCode.EVENT_NOT_PROCESSED, request); + + verify(s3Service, never()).objectExists(IMAGE_KEY); + } + + @Test + void 다른_세션_CCTV_eventId의_object_key를_거부한다() { + given(observationRepository.findByEventId(EVENT_ID.toString())).willReturn(Optional.of(item)); + ConnectEventImageRequest wrongSession = new ConnectEventImageRequest( + "training/" + UUID.randomUUID() + "/events/" + CCTV_CODE + "/" + EVENT_ID + ".jpg", + request.uploadedAt() + ); + ConnectEventImageRequest wrongCctv = new ConnectEventImageRequest( + "training/" + SESSION_ID + "/events/CCTV_002/" + EVENT_ID + ".jpg", + request.uploadedAt() + ); + ConnectEventImageRequest wrongEvent = new ConnectEventImageRequest( + "training/" + SESSION_ID + "/events/" + CCTV_CODE + "/" + UUID.randomUUID() + ".jpg", + request.uploadedAt() + ); + + assertError(CongestionErrorCode.EVENT_IMAGE_IDENTITY_MISMATCH, wrongSession); + assertError(CongestionErrorCode.EVENT_IMAGE_IDENTITY_MISMATCH, wrongCctv); + assertError(CongestionErrorCode.EVENT_IMAGE_IDENTITY_MISMATCH, wrongEvent); + + verify(s3Service, never()).objectExists(IMAGE_KEY); + } + + @Test + void 형식이_잘못된_object_key를_거부한다() { + given(observationRepository.findByEventId(EVENT_ID.toString())).willReturn(Optional.of(item)); + ConnectEventImageRequest malformed = new ConnectEventImageRequest( + "training/" + SESSION_ID + "/monitoring/" + CCTV_CODE + "/" + EVENT_ID + ".jpg", + request.uploadedAt() + ); + + assertError(CongestionErrorCode.EVENT_IMAGE_KEY_INVALID, malformed); + } + + @Test + void S3에_객체가_없으면_연결하지_않는다() { + given(observationRepository.findByEventId(EVENT_ID.toString())).willReturn(Optional.of(item)); + given(s3Service.objectExists(IMAGE_KEY)).willReturn(false); + + assertError(CongestionErrorCode.EVENT_IMAGE_OBJECT_NOT_FOUND, request); + + verify(observationRepository, never()).completeImageUpload( + EVENT_ID.toString(), IMAGE_KEY, request.uploadedAt() + ); + } + + @Test + void 동일한_완료_PATCH는_멱등하게_처리하고_이벤트를_재발행한다() { + item.setEventImageKey(IMAGE_KEY); + item.setImageUploadedAt(request.uploadedAt()); + item.setImageUploadStatus(ImageUploadStatus.COMPLETED); + given(observationRepository.findByEventId(EVENT_ID.toString())).willReturn(Optional.of(item)); + + service.connectImage(principal, EVENT_ID, request); + + verify(s3Service, never()).objectExists(IMAGE_KEY); + verify(observationRepository, never()).completeImageUpload( + EVENT_ID.toString(), IMAGE_KEY, request.uploadedAt() + ); + verify(trainingEventPublisher).publishCongestionImageUpdated(SESSION_ID, item); + } + + @Test + @SuppressWarnings("deprecation") + void 기존_UPLOADED_완료_항목도_멱등하게_처리한다() { + item.setEventImageKey(IMAGE_KEY); + item.setImageUploadedAt(request.uploadedAt()); + item.setImageUploadStatus(ImageUploadStatus.UPLOADED); + given(observationRepository.findByEventId(EVENT_ID.toString())).willReturn(Optional.of(item)); + + service.connectImage(principal, EVENT_ID, request); + + verify(s3Service, never()).objectExists(IMAGE_KEY); + verify(observationRepository, never()).completeImageUpload( + EVENT_ID.toString(), IMAGE_KEY, request.uploadedAt() + ); + verify(trainingEventPublisher).publishCongestionImageUpdated(SESSION_ID, item); + } + + @Test + void 기존_항목의_누락된_이미지_상태는_PENDING으로_취급한다() { + item.setImageUploadStatus(null); + given(observationRepository.findByEventId(EVENT_ID.toString())).willReturn(Optional.of(item)); + given(s3Service.objectExists(IMAGE_KEY)).willReturn(true); + given(observationRepository.completeImageUpload( + EVENT_ID.toString(), IMAGE_KEY, request.uploadedAt() + )).willReturn(true); + + service.connectImage(principal, EVENT_ID, request); + + verify(observationRepository).completeImageUpload( + EVENT_ID.toString(), IMAGE_KEY, request.uploadedAt() + ); + } + + @ParameterizedTest + @EnumSource( + value = DeviceErrorCode.class, + names = {"CCTV_CODE_MISMATCH", "CCTV_DISABLED"} + ) + void CCTV_권한_검증이_실패하면_후속_처리를_중단한다(DeviceErrorCode errorCode) { + given(observationRepository.findByEventId(EVENT_ID.toString())).willReturn(Optional.of(item)); + given(deviceAuthorizationService.validateCctv(principal, CCTV_CODE)) + .willThrow(new ApiException(errorCode)); + + assertThatThrownBy(() -> service.connectImage(principal, EVENT_ID, request)) + .isInstanceOf(ApiException.class) + .hasFieldOrPropertyWithValue("errorCode", errorCode); + + verify(s3Service, never()).objectExists(IMAGE_KEY); + verify(observationRepository, never()).completeImageUpload( + EVENT_ID.toString(), IMAGE_KEY, request.uploadedAt() + ); + verify(trainingEventPublisher, never()).publishCongestionImageUpdated(SESSION_ID, item); + } + + @Test + void 조건부_갱신_경쟁에서_다른_이미지가_완료되면_409를_반환한다() { + ObservationItem latest = item; + latest.setEventImageKey("training/" + SESSION_ID + "/events/" + CCTV_CODE + "/other.jpg"); + latest.setImageUploadedAt(1L); + latest.setImageUploadStatus(ImageUploadStatus.COMPLETED); + ObservationItem initial = ObservationItem.create( + EVENT_ID, SESSION_ID, EDGE_ID, CCTV_CODE, 5.0, 8, 25, 2.5, + CongestionLevel.CROWDED, 1_000L, 2_000L, 2_000L, null, 1L + ); + initial.setEventStatus(EventProcessingStatus.PROCESSED); + given(observationRepository.findByEventId(EVENT_ID.toString())) + .willReturn(Optional.of(initial), Optional.of(latest)); + given(s3Service.objectExists(IMAGE_KEY)).willReturn(true); + given(observationRepository.completeImageUpload( + EVENT_ID.toString(), IMAGE_KEY, request.uploadedAt() + )).willReturn(false); + + assertError(CongestionErrorCode.EVENT_IMAGE_STATE_CONFLICT, request); + } + + private void assertError( + CongestionErrorCode errorCode, + ConnectEventImageRequest connectRequest + ) { + assertThatThrownBy(() -> service.connectImage(principal, EVENT_ID, connectRequest)) + .isInstanceOf(ApiException.class) + .hasFieldOrPropertyWithValue("errorCode", errorCode); + } +} diff --git a/src/test/java/com/saferoute/domain/telemetry/dynamo/entity/TelemetryItemTest.java b/src/test/java/com/saferoute/domain/telemetry/dynamo/entity/TelemetryItemTest.java index ec056587..fdab6e02 100644 --- a/src/test/java/com/saferoute/domain/telemetry/dynamo/entity/TelemetryItemTest.java +++ b/src/test/java/com/saferoute/domain/telemetry/dynamo/entity/TelemetryItemTest.java @@ -40,6 +40,7 @@ class TelemetryItemTest { assertThat(item.getGsi1Pk()).isEqualTo("SESSION#" + SESSION_ID + "#CCTV#CCTV_001"); assertThat(item.getGsi1Sk()).isEqualTo("TIME#1786500005000"); assertThat(item.getEventStatus()).isEqualTo(EventProcessingStatus.RECEIVED); + assertThat(item.getImageUploadStatus()).isEqualTo(ImageUploadStatus.PENDING); } @Test diff --git a/src/test/java/com/saferoute/domain/telemetry/dynamo/repository/ObservationRepositoryTest.java b/src/test/java/com/saferoute/domain/telemetry/dynamo/repository/ObservationRepositoryTest.java index 3044a516..94b09069 100644 --- a/src/test/java/com/saferoute/domain/telemetry/dynamo/repository/ObservationRepositoryTest.java +++ b/src/test/java/com/saferoute/domain/telemetry/dynamo/repository/ObservationRepositoryTest.java @@ -168,6 +168,39 @@ void setUp() { assertThat(captor.getValue().item().getEventStatus().name()).isEqualTo("FAILED"); } + @Test + void PROCESSED이고_이미지가_PENDING_FAILED일_때만_이미지를_완료한다() { + ArgumentCaptor> captor = updateRequestCaptor(); + + boolean completed = repository.completeImageUpload( + "observation-1", "training/session/events/CCTV_001/event.jpg", 2_000L + ); + + assertThat(completed).isTrue(); + verify(table).updateItem(captor.capture()); + UpdateItemEnhancedRequest request = captor.getValue(); + assertThat(request.item().getEventImageKey()) + .isEqualTo("training/session/events/CCTV_001/event.jpg"); + assertThat(request.item().getImageUploadedAt()).isEqualTo(2_000L); + assertThat(request.item().getImageUploadStatus().name()).isEqualTo("COMPLETED"); + assertThat(request.conditionExpression().expression()) + .contains("#eventStatus = :processed") + .contains("attribute_not_exists(#imageStatus)") + .contains("#imageStatus = :pending OR #imageStatus = :failed"); + } + + @Test + void 이미지_조건부_갱신에_실패하면_false를_반환한다() { + doThrow(ConditionalCheckFailedException.builder().message("conflict").build()) + .when(table).updateItem(any(UpdateItemEnhancedRequest.class)); + + boolean completed = repository.completeImageUpload( + "observation-1", "training/session/events/CCTV_001/event.jpg", 2_000L + ); + + assertThat(completed).isFalse(); + } + @SuppressWarnings({"unchecked", "rawtypes"}) private ArgumentCaptor> requestCaptor() { return (ArgumentCaptor) ArgumentCaptor.forClass(PutItemEnhancedRequest.class); diff --git a/src/test/java/com/saferoute/infrastructure/s3/S3ServiceTest.java b/src/test/java/com/saferoute/infrastructure/s3/S3ServiceTest.java index bb326402..12dd3c66 100644 --- a/src/test/java/com/saferoute/infrastructure/s3/S3ServiceTest.java +++ b/src/test/java/com/saferoute/infrastructure/s3/S3ServiceTest.java @@ -22,9 +22,13 @@ import org.mockito.junit.jupiter.MockitoExtension; import org.springframework.mock.web.MockMultipartFile; import software.amazon.awssdk.core.sync.RequestBody; +import software.amazon.awssdk.core.exception.SdkClientException; import software.amazon.awssdk.services.s3.S3Client; import software.amazon.awssdk.services.s3.model.PutObjectRequest; import software.amazon.awssdk.services.s3.model.PutObjectResponse; +import software.amazon.awssdk.services.s3.model.HeadObjectRequest; +import software.amazon.awssdk.services.s3.model.HeadObjectResponse; +import software.amazon.awssdk.services.s3.model.S3Exception; @ExtendWith(MockitoExtension.class) class S3ServiceTest { @@ -92,4 +96,46 @@ void rejectsEmptyFile() { verify(s3Client, never()) .putObject(any(PutObjectRequest.class), any(RequestBody.class)); } + + @Test + void findsExistingObjectWithHeadRequest() { + when(s3Client.headObject(any(HeadObjectRequest.class))) + .thenReturn(HeadObjectResponse.builder().build()); + + assertThat(s3Service.objectExists("training/session/events/CCTV_001/event.jpg")).isTrue(); + + ArgumentCaptor captor = ArgumentCaptor.forClass(HeadObjectRequest.class); + verify(s3Client).headObject(captor.capture()); + assertThat(captor.getValue().bucket()).isEqualTo("test-bucket"); + assertThat(captor.getValue().key()) + .isEqualTo("training/session/events/CCTV_001/event.jpg"); + } + + @Test + void returnsFalseWhenObjectDoesNotExist() { + when(s3Client.headObject(any(HeadObjectRequest.class))) + .thenThrow(S3Exception.builder().statusCode(404).message("not found").build()); + + assertThat(s3Service.objectExists("missing.jpg")).isFalse(); + } + + @Test + void mapsS3ServerErrorToObjectCheckFailed() { + when(s3Client.headObject(any(HeadObjectRequest.class))) + .thenThrow(S3Exception.builder().statusCode(503).message("unavailable").build()); + + assertThatThrownBy(() -> s3Service.objectExists("event.jpg")) + .isInstanceOf(ApiException.class) + .hasFieldOrPropertyWithValue("errorCode", S3ErrorCode.OBJECT_CHECK_FAILED); + } + + @Test + void mapsSdkFailureToObjectCheckFailed() { + when(s3Client.headObject(any(HeadObjectRequest.class))) + .thenThrow(SdkClientException.builder().message("network error").build()); + + assertThatThrownBy(() -> s3Service.objectExists("event.jpg")) + .isInstanceOf(ApiException.class) + .hasFieldOrPropertyWithValue("errorCode", S3ErrorCode.OBJECT_CHECK_FAILED); + } }