From bda3b254ed73fe5213311844b085a16b4f7413f9 Mon Sep 17 00:00:00 2001 From: jmj Date: Wed, 19 Aug 2026 02:54:17 +0900 Subject: [PATCH] =?UTF-8?q?fix(backend):=20=EC=9D=8C=EC=84=B1=20=EB=8B=B5?= =?UTF-8?q?=EB=B3=80=EC=9D=B4=20=EC=A1=B0=EC=9A=A9=ED=9E=88=20=EC=82=AC?= =?UTF-8?q?=EB=9D=BC=EC=A7=88=20=EC=88=98=20=EC=9E=88=EB=8D=98=20=EB=AC=B8?= =?UTF-8?q?=EC=A0=9C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `VoiceAnswerUploadService.submit()` 이 하나의 `@Transactional` 안에서 S3 PUT 과 RabbitMQ 발행까지 전부 처리하고 있었다. 두 가지가 문제였다. **1. 발행이 커밋보다 먼저 나갔다.** publish 이후 커밋이 실패하면 AI 는 존재하지 않는 메시지로 STT 를 돌리고, 돌아온 `callback.voice` 는 `VoiceCallbackService` 에서 "message not found" 로 드롭된다. 사용자는 답변을 말했는데 아무 일도 일어나지 않고 에러도 안 뜬다. Core→AI 발행 중 이 경로만 유일하게 AFTER_COMMIT 이 아니었다 (questions·followup·tts·feedback·analyze.resume 등은 모두 이미 커밋 후 발행). **2. 업로드가 DB 커넥션을 붙잡았다.** 최대 25MB 업로드가 끝날 때까지 트랜잭션이 열려 있어 동시 음성 답변이 커넥션 풀을 잠식한다. ## 변경 - `VoiceAnalysisRequester` 신설 — `VoiceAnswerUploadedEvent` 를 `AFTER_COMMIT` 에 받아 `analyze.voice` 발행. 기존 `SessionTtsRequester`·`SessionFollowupRequester` 와 같은 패턴 - `VoiceAnswerSubmitService` 신설 — 오케스트레이션(검증 → placeholder(tx) → S3 → 부착(tx)). S3 PUT 이 트랜잭션 경계 밖으로 나갔다. `@Transactional` 메서드를 같은 빈에서 자기호출하면 프록시를 타지 않으므로 오케스트레이션을 별 빈으로 분리했다 - `VoiceAnswerUploadService` 는 DB 쓰기 단계만 담당 (`createVoicePlaceholder` / `attachAudioAndRequestAnalysis` / `describe` / `failVoiceUpload`) ## 업로드 실패 보상 S3 PUT 이 트랜잭션 밖으로 나가면서 placeholder 는 이미 커밋된 상태가 된다. 그대로 두면 STT 콜백이 영원히 오지 않아 "음성 인식 중…" 에서 턴이 잠기므로, 업로드 실패 시 메시지를 `FAILED` 로 확정하고 SSE 로 알린다(`VOICE_UPLOAD_FAILED`). 프론트는 기존 STT 실패 경로와 동일하게 턴을 풀고, 사용자는 같은 질문에 다시 답할 수 있다 (`InterviewMessageService.resolveAnswerParent` 의 FAILED 음성 재답변 경로). 부수 효과로 같은 Idempotency-Key 재시도가 실제로 복구 경로가 된다 — 예전에는 전체가 롤백돼 placeholder 자체가 사라졌지만, 이제 남은 메시지에 오디오만 다시 붙인다. ## 테스트 - `VoiceAnalysisRequesterTest` 신설 (2) — 페이로드 필드·메시지 부재 시 스킵 - `VoiceAnswerSubmitServiceTest` 신설 (7) — S3 PUT 이 placeholder 생성 이후·부착 이전에 일어나는지 `InOrder` 로 고정, 업로드 실패 시 FAILED 보상, 멱등 재요청 시 재업로드 안 함, 코덱 파라미터 MIME 키 매핑, 검증 3종 - `VoiceAnswerUploadServiceTest` 재작성 (11) — placeholder 생성 검증 5종 + 부착/보상 `./gradlew cleanTest test` 통과 (ArchUnit 8룰 포함). 음성 관련 25건 전부 green. `docs/messaging.md §5.14` 에 "커밋 후 발행" 규약과 발행 주체 표를 추가했다 — 이 버그가 다시 들어오지 않게 하는 게 목적이다. --- .../common/exception/ApiErrorCode.java | 1 + .../application/VoiceAnalysisRequester.java | 61 ++++ .../application/VoiceAnswerSubmitService.java | 96 ++++++ .../application/VoiceAnswerUploadService.java | 122 +++----- .../event/VoiceAnswerUploadedEvent.java | 13 + .../presentation/VoiceAnswerController.java | 6 +- .../VoiceAnalysisRequesterTest.java | 80 +++++ .../VoiceAnswerSubmitServiceTest.java | 187 ++++++++++++ .../VoiceAnswerUploadServiceTest.java | 287 +++++++++--------- docs/messaging.md | 31 ++ 10 files changed, 657 insertions(+), 227 deletions(-) create mode 100644 backend/src/main/java/com/stackup/stackup/session/application/VoiceAnalysisRequester.java create mode 100644 backend/src/main/java/com/stackup/stackup/session/application/VoiceAnswerSubmitService.java create mode 100644 backend/src/main/java/com/stackup/stackup/session/application/event/VoiceAnswerUploadedEvent.java create mode 100644 backend/src/test/java/com/stackup/stackup/session/application/VoiceAnalysisRequesterTest.java create mode 100644 backend/src/test/java/com/stackup/stackup/session/application/VoiceAnswerSubmitServiceTest.java diff --git a/backend/src/main/java/com/stackup/stackup/common/exception/ApiErrorCode.java b/backend/src/main/java/com/stackup/stackup/common/exception/ApiErrorCode.java index b1611db0..2a50c9e6 100644 --- a/backend/src/main/java/com/stackup/stackup/common/exception/ApiErrorCode.java +++ b/backend/src/main/java/com/stackup/stackup/common/exception/ApiErrorCode.java @@ -48,6 +48,7 @@ public enum ApiErrorCode { VOICE_FILE_TOO_LARGE(HttpStatus.BAD_REQUEST, "음성 파일 크기가 너무 큽니다."), VOICE_INVALID_CONTENT_TYPE(HttpStatus.BAD_REQUEST, "지원하지 않는 음성 형식입니다."), VOICE_MESSAGE_NOT_FOUND(HttpStatus.NOT_FOUND, "음성 메시지를 찾을 수 없습니다."), + VOICE_UPLOAD_FAILED(HttpStatus.INTERNAL_SERVER_ERROR, "음성 파일 업로드에 실패했습니다."), VALIDATION_ERROR(HttpStatus.BAD_REQUEST, "요청 값이 올바르지 않습니다."), ACCESS_DENIED(HttpStatus.FORBIDDEN, "접근 권한이 없습니다."), diff --git a/backend/src/main/java/com/stackup/stackup/session/application/VoiceAnalysisRequester.java b/backend/src/main/java/com/stackup/stackup/session/application/VoiceAnalysisRequester.java new file mode 100644 index 00000000..b66f2e14 --- /dev/null +++ b/backend/src/main/java/com/stackup/stackup/session/application/VoiceAnalysisRequester.java @@ -0,0 +1,61 @@ +package com.stackup.stackup.session.application; + +import com.stackup.stackup.common.config.properties.RabbitMqProperties; +import com.stackup.stackup.common.messaging.MessageContext; +import com.stackup.stackup.common.messaging.RabbitMessagePublisher; +import com.stackup.stackup.session.application.dto.AnalyzeVoicePayload; +import com.stackup.stackup.session.application.event.VoiceAnswerUploadedEvent; +import com.stackup.stackup.session.domain.InterviewMessage; +import com.stackup.stackup.session.domain.InterviewMessageRepository; +import com.stackup.stackup.session.domain.InterviewSession; +import lombok.RequiredArgsConstructor; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Component; +import org.springframework.transaction.annotation.Propagation; +import org.springframework.transaction.annotation.Transactional; +import org.springframework.transaction.event.TransactionPhase; +import org.springframework.transaction.event.TransactionalEventListener; + +// 음성 답변 오디오 부착 commit 후 발화 → analyze.voice envelope 발행 (US-21). +// 텍스트 답변(SessionFollowupRequester) · 질문 TTS(SessionTtsRequester) 와 동일한 AFTER_COMMIT 패턴. +@Component +@RequiredArgsConstructor +public class VoiceAnalysisRequester { + + private static final Logger log = LoggerFactory.getLogger(VoiceAnalysisRequester.class); + + private final RabbitMessagePublisher publisher; + private final RabbitMqProperties properties; + private final InterviewMessageRepository messageRepository; + + @Transactional(readOnly = true, propagation = Propagation.REQUIRES_NEW) + @TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) + public void onVoiceAnswerUploaded(VoiceAnswerUploadedEvent event) { + InterviewMessage message = messageRepository.findById(event.messageId()).orElse(null); + if (message == null) { + log.warn("analyze.voice skipped — message not found. messageId={}", event.messageId()); + return; + } + InterviewSession session = message.getSession(); + InterviewMessage parent = message.getParentMessage(); + + AnalyzeVoicePayload payload = new AnalyzeVoicePayload( + session.getId(), + message.getId(), + parent == null ? null : parent.getId(), + event.audioS3Key(), + event.contentType(), + parent == null ? null : parent.getContent(), + session.getMode().name(), + session.getJobCategory().name() + ); + publisher.publishToAi( + properties.routingKeys().analyzeVoice(), + payload, + new MessageContext(event.userId(), session.getId(), null, null) + ); + log.info("analyze.voice published. sessionId={}, messageId={}, key={}", + session.getId(), message.getId(), event.audioS3Key()); + } +} diff --git a/backend/src/main/java/com/stackup/stackup/session/application/VoiceAnswerSubmitService.java b/backend/src/main/java/com/stackup/stackup/session/application/VoiceAnswerSubmitService.java new file mode 100644 index 00000000..082f1fca --- /dev/null +++ b/backend/src/main/java/com/stackup/stackup/session/application/VoiceAnswerSubmitService.java @@ -0,0 +1,96 @@ +package com.stackup.stackup.session.application; + +import com.stackup.stackup.common.exception.ApiErrorCode; +import com.stackup.stackup.common.exception.DomainException; +import com.stackup.stackup.common.storage.ObjectStorageClient; +import com.stackup.stackup.session.application.VoiceAnswerUploadService.VoicePlaceholder; +import com.stackup.stackup.session.application.dto.MessageResult; +import com.stackup.stackup.session.application.dto.VoiceAnswerUploadCommand; +import java.util.Set; +import lombok.RequiredArgsConstructor; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.springframework.stereotype.Service; + +// 음성 답변 업로드 오케스트레이션: 검증 → placeholder INSERT(tx) → S3 PUT → 오디오 부착(tx). +// +// **트랜잭션 밖에서 S3 에 올린다.** 예전에는 이 메서드 전체가 @Transactional 이었고 그 안에서 +// S3 PUT 과 RabbitMQ 발행을 했다. 두 가지가 문제였다: +// 1) 발행이 commit 보다 먼저 나가서, 이후 commit 이 실패하면 AI 는 존재하지 않는 메시지로 +// STT 를 돌리고 콜백은 "message not found" 로 드롭 — 사용자 답변이 조용히 증발했다. +// 2) 최대 25MB 업로드가 끝날 때까지 DB 커넥션을 붙잡아 커넥션 풀을 잠식했다. +// 발행은 VoiceAnalysisRequester 가 AFTER_COMMIT 에 하고, S3 는 트랜잭션 경계 밖으로 뺐다. +@Service +@RequiredArgsConstructor +public class VoiceAnswerSubmitService { + + private static final Logger log = LoggerFactory.getLogger(VoiceAnswerSubmitService.class); + private static final Set ALLOWED_CONTENT_TYPES = Set.of( + "audio/webm", "audio/ogg", "audio/mpeg", "audio/mp4", "audio/wav", + "audio/x-wav", "audio/m4a", "audio/x-m4a" + ); + private static final long MAX_BYTES = 25L * 1024 * 1024; // Whisper API 25MB 제한 + + private final VoiceAnswerUploadService uploadService; + private final ObjectStorageClient storage; + + public MessageResult submit(Long userId, Long sessionId, VoiceAnswerUploadCommand cmd) { + validate(cmd); + VoicePlaceholder vp = uploadService.createVoicePlaceholder(userId, sessionId, cmd.idempotencyKey()); + Long messageId = vp.placeholder().getId(); + + // idempotency 재호출이면 이미 업로드된 메시지 — 재업로드/재발행 없이 현재 상태를 반환. + if (vp.placeholder().getAudioFilePath() != null) { + return uploadService.describe(messageId); + } + + String key = buildKey(sessionId, messageId, cmd.contentType()); + try { + storage.put(key, cmd.content(), cmd.size(), cmd.contentType()); + } catch (RuntimeException e) { + // placeholder 는 이미 commit 됐다. 그냥 두면 STT 콜백이 오지 않아 턴이 잠기므로 + // FAILED 로 확정해 사용자가 다시 답할 수 있게 한다. + uploadService.failVoiceUpload(sessionId, messageId); + log.warn("voice audio upload to storage failed. sessionId={}, messageId={}, key={}", + sessionId, messageId, key, e); + throw new DomainException(ApiErrorCode.VOICE_UPLOAD_FAILED); + } + return uploadService.attachAudioAndRequestAnalysis( + userId, sessionId, messageId, key, cmd.contentType()); + } + + private void validate(VoiceAnswerUploadCommand cmd) { + if (cmd == null || cmd.content() == null || cmd.size() <= 0) { + throw new DomainException(ApiErrorCode.VOICE_EMPTY_FILE); + } + if (cmd.size() > MAX_BYTES) { + throw new DomainException(ApiErrorCode.VOICE_FILE_TOO_LARGE); + } + if (baseContentType(cmd.contentType()) == null) { + throw new DomainException(ApiErrorCode.VOICE_INVALID_CONTENT_TYPE); + } + } + + // 브라우저 MediaRecorder 는 "audio/webm;codecs=opus" 처럼 코덱 파라미터를 붙인다. + // 파라미터를 떼고 base MIME 만으로 허용 여부를 판단한다. 허용 외면 null. + private static String baseContentType(String contentType) { + if (contentType == null) { + return null; + } + String base = contentType.split(";", 2)[0].trim().toLowerCase(); + return ALLOWED_CONTENT_TYPES.contains(base) ? base : null; + } + + private static String buildKey(Long sessionId, Long messageId, String contentType) { + String base = baseContentType(contentType); + String ext = switch (base == null ? "" : base) { + case "audio/webm" -> "webm"; + case "audio/ogg" -> "ogg"; + case "audio/mpeg" -> "mp3"; + case "audio/mp4", "audio/m4a", "audio/x-m4a" -> "m4a"; + case "audio/wav", "audio/x-wav" -> "wav"; + default -> "bin"; + }; + return "interview/voice/raw/%d/%d.%s".formatted(sessionId, messageId, ext); + } +} diff --git a/backend/src/main/java/com/stackup/stackup/session/application/VoiceAnswerUploadService.java b/backend/src/main/java/com/stackup/stackup/session/application/VoiceAnswerUploadService.java index ce1501e7..68c874b2 100644 --- a/backend/src/main/java/com/stackup/stackup/session/application/VoiceAnswerUploadService.java +++ b/backend/src/main/java/com/stackup/stackup/session/application/VoiceAnswerUploadService.java @@ -1,52 +1,44 @@ package com.stackup.stackup.session.application; -import com.stackup.stackup.common.config.properties.RabbitMqProperties; import com.stackup.stackup.common.exception.ApiErrorCode; import com.stackup.stackup.common.exception.DomainException; -import com.stackup.stackup.common.messaging.MessageContext; -import com.stackup.stackup.common.messaging.RabbitMessagePublisher; -import com.stackup.stackup.common.storage.ObjectStorageClient; -import com.stackup.stackup.session.application.dto.AnalyzeVoicePayload; +import com.stackup.stackup.common.messaging.RealtimeNotifyEvent; +import com.stackup.stackup.common.sse.SseEventType; import com.stackup.stackup.session.application.dto.MessageResult; -import com.stackup.stackup.session.application.dto.VoiceAnswerUploadCommand; +import com.stackup.stackup.session.application.event.VoiceAnswerUploadedEvent; import com.stackup.stackup.session.domain.InterviewMessage; import com.stackup.stackup.session.domain.InterviewMessageRepository; import com.stackup.stackup.session.domain.InterviewSession; import com.stackup.stackup.session.domain.InterviewSessionRepository; import com.stackup.stackup.session.domain.MessageRole; import com.stackup.stackup.session.domain.SessionStatus; -import java.util.Set; import lombok.RequiredArgsConstructor; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.springframework.context.ApplicationEventPublisher; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; -// 음성 답변 업로드: S3 PUT → InterviewMessage INSERT (content=null, audio_file_path=key) → analyze.voice 발행. -// STT/분석 결과는 callback.voice 로 도착해서 InterviewMessage.content 채움 + voice metrics INSERT + followup 트리거. +// 음성 답변의 DB 쓰기 단계. S3 PUT · RabbitMQ 발행은 여기서 하지 않는다. +// - 업로드 오케스트레이션(검증 → S3 → 이 서비스 호출): VoiceAnswerSubmitService +// - analyze.voice 발행: VoiceAnalysisRequester (AFTER_COMMIT) +// 스트리밍 음성(stream-begin)은 createVoicePlaceholder 만 쓴다 — 오디오는 WS 로 흐른다. @Service @RequiredArgsConstructor public class VoiceAnswerUploadService { private static final Logger log = LoggerFactory.getLogger(VoiceAnswerUploadService.class); - private static final Set ALLOWED_CONTENT_TYPES = Set.of( - "audio/webm", "audio/ogg", "audio/mpeg", "audio/mp4", "audio/wav", - "audio/x-wav", "audio/m4a", "audio/x-m4a" - ); - private static final long MAX_BYTES = 25L * 1024 * 1024; // Whisper API 25MB 제한 private final InterviewSessionRepository sessionRepository; private final InterviewMessageRepository messageRepository; - private final ObjectStorageClient storage; - private final RabbitMessagePublisher publisher; - private final RabbitMqProperties properties; + private final ApplicationEventPublisher events; // 스트리밍/배치 음성 답변이 공유하는 placeholder 생성 결과. public record VoicePlaceholder(InterviewSession session, InterviewMessage placeholder, InterviewMessage parentQuestion) {} // 세션 조회 → idempotency → 상태/직전메시지 검증 → placeholder save. - // 배치 업로드(submit)와 스트리밍 시작(VoiceStreamService) 양쪽에서 재사용한다. + // 배치 업로드(VoiceAnswerSubmitService)와 스트리밍 시작(VoiceStreamService) 양쪽에서 재사용한다. @Transactional public VoicePlaceholder createVoicePlaceholder(Long userId, Long sessionId, String idempotencyKey) { InterviewSession session = sessionRepository.findByIdAndUser_IdAndDeletedFalse(sessionId, userId) @@ -76,71 +68,47 @@ public VoicePlaceholder createVoicePlaceholder(Long userId, Long sessionId, Stri return new VoicePlaceholder(session, placeholder, latest); } + // 업로드된 오디오 키를 메시지에 붙이고 analyze.voice 요청 이벤트를 발행한다. + // 발행은 AFTER_COMMIT 리스너가 받으므로, 이 트랜잭션이 롤백되면 AI 호출도 일어나지 않는다. @Transactional - public MessageResult submit(Long userId, Long sessionId, VoiceAnswerUploadCommand cmd) { - validate(cmd); - VoicePlaceholder vp = createVoicePlaceholder(userId, sessionId, cmd.idempotencyKey()); - // idempotency 재호출이면 이미 완료된 메시지일 수 있음 — 기존 동작 유지: 그대로 반환. - InterviewMessage placeholder = vp.placeholder(); - if (placeholder.getAudioFilePath() != null) { - return MessageResult.of(placeholder); // 이미 업로드됨(중복) - } - String key = buildKey(sessionId, placeholder.getId(), cmd.contentType()); - storage.put(key, cmd.content(), cmd.size(), cmd.contentType()); - placeholder.attachAudio(key); - - AnalyzeVoicePayload payload = new AnalyzeVoicePayload( - sessionId, - placeholder.getId(), - vp.parentQuestion().getId(), - key, - cmd.contentType(), - vp.parentQuestion().getContent(), - vp.session().getMode().name(), - vp.session().getJobCategory().name() - ); - publisher.publishToAi( - properties.routingKeys().analyzeVoice(), - payload, - new MessageContext(userId, sessionId, null, null) - ); - log.info("analyze.voice published. sessionId={}, messageId={}, key={}", - sessionId, placeholder.getId(), key); - return MessageResult.of(placeholder); - } + public MessageResult attachAudioAndRequestAnalysis( + Long userId, Long sessionId, Long messageId, String audioS3Key, String contentType) { + InterviewMessage message = messageRepository.findById(messageId) + .orElseThrow(() -> new DomainException(ApiErrorCode.VOICE_MESSAGE_NOT_FOUND)); - private void validate(VoiceAnswerUploadCommand cmd) { - if (cmd == null || cmd.content() == null || cmd.size() <= 0) { - throw new DomainException(ApiErrorCode.VOICE_EMPTY_FILE); - } - if (cmd.size() > MAX_BYTES) { - throw new DomainException(ApiErrorCode.VOICE_FILE_TOO_LARGE); - } - if (baseContentType(cmd.contentType()) == null) { - throw new DomainException(ApiErrorCode.VOICE_INVALID_CONTENT_TYPE); + // 같은 Idempotency-Key 재요청이 경합해 이미 붙었으면 재발행하지 않는다. + if (message.getAudioFilePath() != null) { + return MessageResult.of(message); } + message.attachAudio(audioS3Key); + events.publishEvent(new VoiceAnswerUploadedEvent( + userId, sessionId, message.getId(), audioS3Key, contentType)); + return MessageResult.of(message); } - // 브라우저 MediaRecorder 는 "audio/webm;codecs=opus" 처럼 코덱 파라미터를 붙인다. - // 파라미터를 떼고 base MIME 만으로 허용 여부를 판단한다. 허용 외면 null. - private static String baseContentType(String contentType) { - if (contentType == null) { - return null; - } - String base = contentType.split(";", 2)[0].trim().toLowerCase(); - return ALLOWED_CONTENT_TYPES.contains(base) ? base : null; + @Transactional(readOnly = true) + public MessageResult describe(Long messageId) { + return MessageResult.of(messageRepository.findById(messageId) + .orElseThrow(() -> new DomainException(ApiErrorCode.VOICE_MESSAGE_NOT_FOUND))); } - private static String buildKey(Long sessionId, Long messageId, String contentType) { - String base = baseContentType(contentType); - String ext = switch (base == null ? "" : base) { - case "audio/webm" -> "webm"; - case "audio/ogg" -> "ogg"; - case "audio/mpeg" -> "mp3"; - case "audio/mp4", "audio/m4a", "audio/x-m4a" -> "m4a"; - case "audio/wav", "audio/x-wav" -> "wav"; - default -> "bin"; - }; - return "interview/voice/raw/%d/%d.%s".formatted(sessionId, messageId, ext); + // S3 업로드가 실패했을 때의 보상. placeholder 를 지우지 않고 FAILED 로 확정한다 — + // 그냥 두면 STT 콜백이 영원히 오지 않아 "음성 인식 중…" 에서 턴이 잠긴다. + // FAILED 로 두면 프론트가 턴을 풀고, 사용자는 같은 질문에 다시 답할 수 있다 + // (InterviewMessageService.resolveAnswerParent 의 FAILED 음성 답변 재답변 경로). + @Transactional + public void failVoiceUpload(Long sessionId, Long messageId) { + InterviewMessage message = messageRepository.findById(messageId).orElse(null); + if (message == null) { + return; + } + message.failVoiceTranscription(); + VoiceCallbackService.VoiceFailedNotice notice = new VoiceCallbackService.VoiceFailedNotice( + sessionId, message.getId(), "VOICE_UPLOAD_FAILED", message.getContent()); + events.publishEvent(RealtimeNotifyEvent.session(sessionId, SseEventType.SESSION_MESSAGE, notice)); + events.publishEvent(RealtimeNotifyEvent.user(message.getSession().getUser().getId(), + SseEventType.SESSION_MESSAGE, notice)); + log.warn("voice upload failed — message marked FAILED. sessionId={}, messageId={}", + sessionId, messageId); } } diff --git a/backend/src/main/java/com/stackup/stackup/session/application/event/VoiceAnswerUploadedEvent.java b/backend/src/main/java/com/stackup/stackup/session/application/event/VoiceAnswerUploadedEvent.java new file mode 100644 index 00000000..7c81af39 --- /dev/null +++ b/backend/src/main/java/com/stackup/stackup/session/application/event/VoiceAnswerUploadedEvent.java @@ -0,0 +1,13 @@ +package com.stackup.stackup.session.application.event; + +// 음성 답변 오디오 키가 메시지에 붙은 뒤 발행. VoiceAnalysisRequester 가 AFTER_COMMIT 에 받아 +// analyze.voice 를 발행한다. commit 전에 발행하면 롤백 시 AI 가 존재하지 않는 메시지로 STT 를 +// 돌리고, 콜백은 "message not found" 로 드롭되어 사용자 답변이 조용히 증발한다. +public record VoiceAnswerUploadedEvent( + Long userId, + Long sessionId, + Long messageId, + String audioS3Key, + String contentType +) { +} diff --git a/backend/src/main/java/com/stackup/stackup/session/presentation/VoiceAnswerController.java b/backend/src/main/java/com/stackup/stackup/session/presentation/VoiceAnswerController.java index 75d7a2b6..d6fb30d7 100644 --- a/backend/src/main/java/com/stackup/stackup/session/presentation/VoiceAnswerController.java +++ b/backend/src/main/java/com/stackup/stackup/session/presentation/VoiceAnswerController.java @@ -1,7 +1,7 @@ package com.stackup.stackup.session.presentation; import com.stackup.stackup.common.security.UserPrincipal; -import com.stackup.stackup.session.application.VoiceAnswerUploadService; +import com.stackup.stackup.session.application.VoiceAnswerSubmitService; import com.stackup.stackup.session.application.VoiceStreamService; import com.stackup.stackup.session.application.dto.MessageResult; import com.stackup.stackup.session.application.dto.VoiceAnswerUploadCommand; @@ -30,7 +30,7 @@ @RequiredArgsConstructor public class VoiceAnswerController { - private final VoiceAnswerUploadService uploadService; + private final VoiceAnswerSubmitService submitService; private final VoiceStreamService streamService; @Operation( @@ -61,7 +61,7 @@ public MessageResponse submit( audio.getOriginalFilename(), idempotencyKey ); - MessageResult result = uploadService.submit(principal.userId(), sessionId, cmd); + MessageResult result = submitService.submit(principal.userId(), sessionId, cmd); return MessageResponse.from(result); } diff --git a/backend/src/test/java/com/stackup/stackup/session/application/VoiceAnalysisRequesterTest.java b/backend/src/test/java/com/stackup/stackup/session/application/VoiceAnalysisRequesterTest.java new file mode 100644 index 00000000..6b36d977 --- /dev/null +++ b/backend/src/test/java/com/stackup/stackup/session/application/VoiceAnalysisRequesterTest.java @@ -0,0 +1,80 @@ +package com.stackup.stackup.session.application; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.stackup.stackup.common.config.properties.RabbitMqProperties; +import com.stackup.stackup.common.messaging.MessageContext; +import com.stackup.stackup.common.messaging.RabbitMessagePublisher; +import com.stackup.stackup.session.application.dto.AnalyzeVoicePayload; +import com.stackup.stackup.session.application.event.VoiceAnswerUploadedEvent; +import com.stackup.stackup.session.domain.InterviewMessage; +import com.stackup.stackup.session.domain.InterviewMessageRepository; +import com.stackup.stackup.session.domain.InterviewSession; +import com.stackup.stackup.session.domain.JobCategory; +import com.stackup.stackup.session.domain.SessionMode; +import java.util.Optional; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.ArgumentCaptor; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.Mockito; +import org.mockito.junit.jupiter.MockitoExtension; + +@ExtendWith(MockitoExtension.class) +class VoiceAnalysisRequesterTest { + + @Mock RabbitMessagePublisher publisher; + @Mock RabbitMqProperties properties; + @Mock RabbitMqProperties.RoutingKeyProperties routingKeys; + @Mock InterviewMessageRepository messageRepository; + @InjectMocks VoiceAnalysisRequester requester; + + @Test + void publishesAnalyzeVoice_afterCommit() { + InterviewSession session = Mockito.mock(InterviewSession.class); + when(session.getId()).thenReturn(10L); + when(session.getMode()).thenReturn(SessionMode.TECHNICAL); + when(session.getJobCategory()).thenReturn(JobCategory.BACKEND); + InterviewMessage parent = Mockito.mock(InterviewMessage.class); + when(parent.getId()).thenReturn(100L); + when(parent.getContent()).thenReturn("Tell me about ACID."); + InterviewMessage message = Mockito.mock(InterviewMessage.class); + when(message.getId()).thenReturn(200L); + when(message.getSession()).thenReturn(session); + when(message.getParentMessage()).thenReturn(parent); + when(messageRepository.findById(200L)).thenReturn(Optional.of(message)); + when(properties.routingKeys()).thenReturn(routingKeys); + when(routingKeys.analyzeVoice()).thenReturn("analyze.voice"); + + requester.onVoiceAnswerUploaded(new VoiceAnswerUploadedEvent( + 1L, 10L, 200L, "interview/voice/raw/10/200.webm", "audio/webm")); + + ArgumentCaptor captor = ArgumentCaptor.forClass(AnalyzeVoicePayload.class); + verify(publisher).publishToAi(eq("analyze.voice"), captor.capture(), any(MessageContext.class)); + AnalyzeVoicePayload payload = captor.getValue(); + assertThat(payload.sessionId()).isEqualTo(10L); + assertThat(payload.messageId()).isEqualTo(200L); + assertThat(payload.parentQuestionMessageId()).isEqualTo(100L); + assertThat(payload.audioS3Key()).isEqualTo("interview/voice/raw/10/200.webm"); + assertThat(payload.contentType()).isEqualTo("audio/webm"); + assertThat(payload.previousQuestionText()).isEqualTo("Tell me about ACID."); + assertThat(payload.mode()).isEqualTo("TECHNICAL"); + assertThat(payload.jobCategory()).isEqualTo("BACKEND"); + } + + @Test + void skips_whenMessageNotFound() { + when(messageRepository.findById(999L)).thenReturn(Optional.empty()); + + requester.onVoiceAnswerUploaded(new VoiceAnswerUploadedEvent( + 1L, 10L, 999L, "k", "audio/webm")); + + verify(publisher, never()).publishToAi(any(), any(), any()); + } +} diff --git a/backend/src/test/java/com/stackup/stackup/session/application/VoiceAnswerSubmitServiceTest.java b/backend/src/test/java/com/stackup/stackup/session/application/VoiceAnswerSubmitServiceTest.java new file mode 100644 index 00000000..99e76599 --- /dev/null +++ b/backend/src/test/java/com/stackup/stackup/session/application/VoiceAnswerSubmitServiceTest.java @@ -0,0 +1,187 @@ +package com.stackup.stackup.session.application; + +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.eq; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoInteractions; +import static org.mockito.Mockito.when; + +import com.stackup.stackup.common.exception.ApiErrorCode; +import com.stackup.stackup.common.exception.DomainException; +import com.stackup.stackup.common.storage.ObjectStorageClient; +import com.stackup.stackup.session.application.VoiceAnswerUploadService.VoicePlaceholder; +import com.stackup.stackup.session.application.dto.MessageResult; +import com.stackup.stackup.session.application.dto.VoiceAnswerUploadCommand; +import com.stackup.stackup.session.domain.InterviewMessage; +import com.stackup.stackup.session.domain.InterviewSession; +import com.stackup.stackup.session.domain.JobCategory; +import com.stackup.stackup.session.domain.SessionMode; +import com.stackup.stackup.user.domain.User; +import java.io.ByteArrayInputStream; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.InOrder; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.test.util.ReflectionTestUtils; + +@ExtendWith(MockitoExtension.class) +class VoiceAnswerSubmitServiceTest { + + private static final String KEY = "interview/voice/raw/10/200.webm"; + + @Mock VoiceAnswerUploadService uploadService; + @Mock ObjectStorageClient storage; + @InjectMocks VoiceAnswerSubmitService service; + + // S3 PUT 은 반드시 placeholder 생성 이후 · 오디오 부착 이전에, 트랜잭션 밖에서 일어난다. + @Test + void submit_putsAudioBetweenPlaceholderAndAttach() { + InterviewSession session = sessionInProgress(10L); + InterviewMessage question = InterviewMessage.interviewer(session, 1, "Tell me about ACID."); + ReflectionTestUtils.setField(question, "id", 100L); + InterviewMessage placeholder = InterviewMessage.voiceInterviewee(session, 2, question, "idem-1"); + ReflectionTestUtils.setField(placeholder, "id", 200L); + MessageResult attached = MessageResult.of(placeholder); + + when(uploadService.createVoicePlaceholder(1L, 10L, "idem-1")) + .thenReturn(new VoicePlaceholder(session, placeholder, question)); + when(uploadService.attachAudioAndRequestAnalysis(1L, 10L, 200L, KEY, "audio/webm")) + .thenReturn(attached); + + MessageResult result = service.submit(1L, 10L, + command("voice-bytes".getBytes(), "audio/webm", "idem-1")); + + assertThat(result).isSameAs(attached); + verify(storage).put(eq(KEY), any(), eq(11L), eq("audio/webm")); + + InOrder order = inOrder(uploadService, storage); + order.verify(uploadService).createVoicePlaceholder(1L, 10L, "idem-1"); + order.verify(storage).put(eq(KEY), any(), anyLong(), any()); + order.verify(uploadService).attachAudioAndRequestAnalysis(1L, 10L, 200L, KEY, "audio/webm"); + } + + // 코덱 파라미터가 붙은 MediaRecorder MIME 도 확장자 매핑이 되어야 한다. + @Test + void submit_stripsCodecParameterWhenBuildingKey() { + InterviewSession session = sessionInProgress(10L); + InterviewMessage question = InterviewMessage.interviewer(session, 1, "q"); + InterviewMessage placeholder = InterviewMessage.voiceInterviewee(session, 2, question, null); + ReflectionTestUtils.setField(placeholder, "id", 200L); + + when(uploadService.createVoicePlaceholder(1L, 10L, null)) + .thenReturn(new VoicePlaceholder(session, placeholder, question)); + when(uploadService.attachAudioAndRequestAnalysis( + 1L, 10L, 200L, KEY, "audio/webm;codecs=opus")).thenReturn(MessageResult.of(placeholder)); + + service.submit(1L, 10L, command("x".getBytes(), "audio/webm;codecs=opus", null)); + + verify(storage).put(eq(KEY), any(), anyLong(), eq("audio/webm;codecs=opus")); + } + + // 같은 Idempotency-Key 재요청: 이미 업로드된 메시지는 재업로드·재발행하지 않는다. + @Test + void submit_returnsExisting_whenAudioAlreadyAttached() { + InterviewSession session = sessionInProgress(10L); + InterviewMessage question = InterviewMessage.interviewer(session, 1, "q"); + InterviewMessage placeholder = InterviewMessage.voiceInterviewee(session, 2, question, "idem-1"); + ReflectionTestUtils.setField(placeholder, "id", 200L); + placeholder.attachAudio(KEY); + MessageResult described = MessageResult.of(placeholder); + + when(uploadService.createVoicePlaceholder(1L, 10L, "idem-1")) + .thenReturn(new VoicePlaceholder(session, placeholder, question)); + when(uploadService.describe(200L)).thenReturn(described); + + MessageResult result = service.submit(1L, 10L, + command("voice".getBytes(), "audio/webm", "idem-1")); + + assertThat(result).isSameAs(described); + verify(storage, never()).put(any(), any(), anyLong(), any()); + verify(uploadService, never()).attachAudioAndRequestAnalysis(any(), any(), any(), any(), any()); + } + + // S3 업로드가 터지면 placeholder 는 이미 commit 돼 있다. FAILED 로 확정해 턴 잠김을 막는다. + @Test + void submit_marksMessageFailed_whenStoragePutThrows() { + InterviewSession session = sessionInProgress(10L); + InterviewMessage question = InterviewMessage.interviewer(session, 1, "q"); + InterviewMessage placeholder = InterviewMessage.voiceInterviewee(session, 2, question, null); + ReflectionTestUtils.setField(placeholder, "id", 200L); + + when(uploadService.createVoicePlaceholder(1L, 10L, null)) + .thenReturn(new VoicePlaceholder(session, placeholder, question)); + doThrow(new RuntimeException("s3 down")) + .when(storage).put(any(), any(), anyLong(), any()); + + assertThatThrownBy(() -> service.submit(1L, 10L, command("v".getBytes(), "audio/webm", null))) + .isInstanceOfSatisfying(DomainException.class, e -> + assertThat(e.getErrorCode()).isEqualTo(ApiErrorCode.VOICE_UPLOAD_FAILED)); + + verify(uploadService).failVoiceUpload(10L, 200L); + verify(uploadService, never()).attachAudioAndRequestAnalysis(any(), any(), any(), any(), any()); + } + + @Test + void submit_rejectsEmptyAudio() { + assertThatThrownBy(() -> service.submit(1L, 10L, command(new byte[0], "audio/webm", null))) + .isInstanceOfSatisfying(DomainException.class, e -> + assertThat(e.getErrorCode()).isEqualTo(ApiErrorCode.VOICE_EMPTY_FILE)); + + verifyNoInteractions(uploadService, storage); + } + + @Test + void submit_rejectsTooLargeAudio() { + VoiceAnswerUploadCommand cmd = new VoiceAnswerUploadCommand( + new ByteArrayInputStream(new byte[] {1}), + 25L * 1024 * 1024 + 1, + "audio/webm", + "a.webm", + null + ); + + assertThatThrownBy(() -> service.submit(1L, 10L, cmd)) + .isInstanceOfSatisfying(DomainException.class, e -> + assertThat(e.getErrorCode()).isEqualTo(ApiErrorCode.VOICE_FILE_TOO_LARGE)); + + verifyNoInteractions(uploadService, storage); + } + + @Test + void submit_rejectsUnsupportedContentType() { + assertThatThrownBy(() -> service.submit(1L, 10L, command("x".getBytes(), "text/plain", null))) + .isInstanceOfSatisfying(DomainException.class, e -> + assertThat(e.getErrorCode()).isEqualTo(ApiErrorCode.VOICE_INVALID_CONTENT_TYPE)); + + verifyNoInteractions(uploadService, storage); + } + + private VoiceAnswerUploadCommand command(byte[] content, String contentType, String idempotencyKey) { + return new VoiceAnswerUploadCommand( + new ByteArrayInputStream(content), + content.length, + contentType, + "answer.webm", + idempotencyKey + ); + } + + private InterviewSession sessionInProgress(Long id) { + User user = User.createGithubUser(1L, "u", null, null, "t"); + ReflectionTestUtils.setField(user, "id", 1L); + InterviewSession s = InterviewSession.create( + user, "t", null, SessionMode.TECHNICAL, JobCategory.BACKEND, 5, 30, null, null + ); + ReflectionTestUtils.setField(s, "id", id); + s.start(); + return s; + } +} diff --git a/backend/src/test/java/com/stackup/stackup/session/application/VoiceAnswerUploadServiceTest.java b/backend/src/test/java/com/stackup/stackup/session/application/VoiceAnswerUploadServiceTest.java index df0a1360..b24f2e13 100644 --- a/backend/src/test/java/com/stackup/stackup/session/application/VoiceAnswerUploadServiceTest.java +++ b/backend/src/test/java/com/stackup/stackup/session/application/VoiceAnswerUploadServiceTest.java @@ -3,21 +3,16 @@ 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.eq; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; -import com.stackup.stackup.common.config.properties.RabbitMqProperties; import com.stackup.stackup.common.exception.ApiErrorCode; import com.stackup.stackup.common.exception.DomainException; -import com.stackup.stackup.common.messaging.MessageContext; -import com.stackup.stackup.common.messaging.RabbitMessagePublisher; -import com.stackup.stackup.common.storage.ObjectStorageClient; -import com.stackup.stackup.session.application.dto.AnalyzeVoicePayload; +import com.stackup.stackup.common.messaging.RealtimeNotifyEvent; +import com.stackup.stackup.session.application.VoiceAnswerUploadService.VoicePlaceholder; import com.stackup.stackup.session.application.dto.MessageResult; -import com.stackup.stackup.session.application.dto.VoiceAnswerUploadCommand; +import com.stackup.stackup.session.application.event.VoiceAnswerUploadedEvent; import com.stackup.stackup.session.domain.InterviewMessage; import com.stackup.stackup.session.domain.InterviewMessageRepository; import com.stackup.stackup.session.domain.InterviewSession; @@ -26,159 +21,204 @@ import com.stackup.stackup.session.domain.MessageStatus; import com.stackup.stackup.session.domain.SessionMode; import com.stackup.stackup.user.domain.User; -import java.io.ByteArrayInputStream; -import java.time.Duration; import java.util.Optional; -import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.extension.ExtendWith; import org.mockito.ArgumentCaptor; +import org.mockito.InjectMocks; import org.mockito.Mock; import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.context.ApplicationEventPublisher; import org.springframework.test.util.ReflectionTestUtils; @ExtendWith(MockitoExtension.class) class VoiceAnswerUploadServiceTest { + private static final String KEY = "interview/voice/raw/10/200.webm"; + @Mock InterviewSessionRepository sessionRepository; @Mock InterviewMessageRepository messageRepository; - @Mock ObjectStorageClient storage; - @Mock RabbitMessagePublisher publisher; - - VoiceAnswerUploadService service; - - @BeforeEach - void setUp() { - service = new VoiceAnswerUploadService( - sessionRepository, - messageRepository, - storage, - publisher, - rabbitMqProperties() - ); - } + @Mock ApplicationEventPublisher events; + @InjectMocks VoiceAnswerUploadService service; + + // ── createVoicePlaceholder ──────────────────────────────────────────────── @Test - void submit_savesAudioAndPublishesAnalyzeVoice() { + void createVoicePlaceholder_insertsAfterLatestQuestion() { InterviewSession session = sessionInProgress(10L); - InterviewMessage question = InterviewMessage.interviewer(session, 1, "Tell me about ACID."); + InterviewMessage question = InterviewMessage.interviewer(session, 3, "Tell me about ACID."); ReflectionTestUtils.setField(question, "id", 100L); when(sessionRepository.findByIdAndUser_IdAndDeletedFalse(10L, 1L)).thenReturn(Optional.of(session)); - when(messageRepository.findFirstBySession_IdOrderBySequenceNumberDesc(10L)).thenReturn(Optional.of(question)); + when(messageRepository.findFirstBySession_IdOrderBySequenceNumberDesc(10L)) + .thenReturn(Optional.of(question)); when(messageRepository.save(any(InterviewMessage.class))).thenAnswer(inv -> { InterviewMessage m = inv.getArgument(0); ReflectionTestUtils.setField(m, "id", 200L); return m; }); - MessageResult result = service.submit(1L, 10L, - command("voice-bytes".getBytes(), "audio/webm", "idem-1")); + VoicePlaceholder vp = service.createVoicePlaceholder(1L, 10L, "idem-1"); - assertThat(result.id()).isEqualTo(200L); - assertThat(result.status()).isEqualTo(MessageStatus.CREATED); - assertThat(result.audioFilePath()).isEqualTo("interview/voice/raw/10/200.webm"); - verify(storage).put(eq("interview/voice/raw/10/200.webm"), any(), eq(11L), eq("audio/webm")); - - ArgumentCaptor payloadCaptor = ArgumentCaptor.forClass(AnalyzeVoicePayload.class); - verify(publisher).publishToAi(eq("analyze.voice"), payloadCaptor.capture(), any(MessageContext.class)); - AnalyzeVoicePayload payload = payloadCaptor.getValue(); - assertThat(payload.sessionId()).isEqualTo(10L); - assertThat(payload.messageId()).isEqualTo(200L); - assertThat(payload.parentQuestionMessageId()).isEqualTo(100L); - assertThat(payload.audioS3Key()).isEqualTo("interview/voice/raw/10/200.webm"); - assertThat(payload.previousQuestionText()).isEqualTo("Tell me about ACID."); + assertThat(vp.placeholder().getId()).isEqualTo(200L); + assertThat(vp.placeholder().getSequenceNumber()).isEqualTo(4); + assertThat(vp.placeholder().getStatus()).isEqualTo(MessageStatus.CREATED); + assertThat(vp.placeholder().getAudioFilePath()).isNull(); + assertThat(vp.parentQuestion()).isSameAs(question); } @Test - void submit_rejectsEmptyAudio() { - assertThatThrownBy(() -> service.submit(1L, 10L, - command(new byte[0], "audio/webm", null))) - .isInstanceOfSatisfying(DomainException.class, exception -> - assertThat(exception.getErrorCode()).isEqualTo(ApiErrorCode.VOICE_EMPTY_FILE)); + void createVoicePlaceholder_returnsExistingOnIdempotencyHit() { + InterviewSession session = sessionInProgress(10L); + InterviewMessage question = InterviewMessage.interviewer(session, 1, "q"); + InterviewMessage existing = InterviewMessage.voiceInterviewee(session, 2, question, "idem-1"); + ReflectionTestUtils.setField(existing, "id", 200L); - verify(sessionRepository, never()).findByIdAndUser_IdAndDeletedFalse(any(), any()); - } + when(sessionRepository.findByIdAndUser_IdAndDeletedFalse(10L, 1L)).thenReturn(Optional.of(session)); + when(messageRepository.findBySession_IdAndIdempotencyKey(10L, "idem-1")) + .thenReturn(Optional.of(existing)); - @Test - void submit_rejectsTooLargeAudio() { - VoiceAnswerUploadCommand cmd = new VoiceAnswerUploadCommand( - new ByteArrayInputStream(new byte[] {1}), - 25L * 1024 * 1024 + 1, - "audio/webm", - "a.webm", - null - ); + VoicePlaceholder vp = service.createVoicePlaceholder(1L, 10L, "idem-1"); - assertThatThrownBy(() -> service.submit(1L, 10L, cmd)) - .isInstanceOfSatisfying(DomainException.class, exception -> - assertThat(exception.getErrorCode()).isEqualTo(ApiErrorCode.VOICE_FILE_TOO_LARGE)); + assertThat(vp.placeholder()).isSameAs(existing); + verify(messageRepository, never()).save(any()); } @Test - void submit_rejectsUnsupportedContentType() { - assertThatThrownBy(() -> service.submit(1L, 10L, - command("x".getBytes(), "text/plain", null))) - .isInstanceOfSatisfying(DomainException.class, exception -> - assertThat(exception.getErrorCode()).isEqualTo(ApiErrorCode.VOICE_INVALID_CONTENT_TYPE)); + void createVoicePlaceholder_rejectsWhenSessionNotFound() { + when(sessionRepository.findByIdAndUser_IdAndDeletedFalse(10L, 1L)).thenReturn(Optional.empty()); + + assertThatThrownBy(() -> service.createVoicePlaceholder(1L, 10L, null)) + .isInstanceOfSatisfying(DomainException.class, e -> + assertThat(e.getErrorCode()).isEqualTo(ApiErrorCode.SESSION_NOT_FOUND)); } @Test - void submit_rejectsWhenSessionIsNotInProgress() { + void createVoicePlaceholder_rejectsWhenSessionIsNotInProgress() { InterviewSession session = sessionFixture(10L); when(sessionRepository.findByIdAndUser_IdAndDeletedFalse(10L, 1L)).thenReturn(Optional.of(session)); - assertThatThrownBy(() -> service.submit(1L, 10L, - command("voice".getBytes(), "audio/webm", null))) - .isInstanceOfSatisfying(DomainException.class, exception -> - assertThat(exception.getErrorCode()).isEqualTo(ApiErrorCode.SESSION_INVALID_STATE)); + assertThatThrownBy(() -> service.createVoicePlaceholder(1L, 10L, null)) + .isInstanceOfSatisfying(DomainException.class, e -> + assertThat(e.getErrorCode()).isEqualTo(ApiErrorCode.SESSION_INVALID_STATE)); - verify(storage, never()).put(any(), any(), anyLong(), any()); - verify(publisher, never()).publishToAi(any(), any(), any()); + verify(messageRepository, never()).save(any()); } @Test - void submit_rejectsWhenNoQuestionMessageExists() { + void createVoicePlaceholder_rejectsWhenNoQuestionMessageExists() { InterviewSession session = sessionInProgress(10L); when(sessionRepository.findByIdAndUser_IdAndDeletedFalse(10L, 1L)).thenReturn(Optional.of(session)); - when(messageRepository.findFirstBySession_IdOrderBySequenceNumberDesc(10L)).thenReturn(Optional.empty()); - - assertThatThrownBy(() -> service.submit(1L, 10L, - command("voice".getBytes(), "audio/webm", null))) - .isInstanceOfSatisfying(DomainException.class, exception -> - assertThat(exception.getErrorCode()).isEqualTo(ApiErrorCode.SESSION_INVALID_STATE)); + when(messageRepository.findFirstBySession_IdOrderBySequenceNumberDesc(10L)) + .thenReturn(Optional.empty()); - verify(storage, never()).put(any(), any(), anyLong(), any()); - verify(publisher, never()).publishToAi(any(), any(), any()); + assertThatThrownBy(() -> service.createVoicePlaceholder(1L, 10L, null)) + .isInstanceOfSatisfying(DomainException.class, e -> + assertThat(e.getErrorCode()).isEqualTo(ApiErrorCode.SESSION_INVALID_STATE)); } @Test - void submit_rejectsWhenLastMessageIsNotQuestion() { + void createVoicePlaceholder_rejectsWhenLastMessageIsNotQuestion() { InterviewSession session = sessionInProgress(10L); InterviewMessage priorAnswer = InterviewMessage.interviewee(session, 1, "already answered", null, null); when(sessionRepository.findByIdAndUser_IdAndDeletedFalse(10L, 1L)).thenReturn(Optional.of(session)); - when(messageRepository.findFirstBySession_IdOrderBySequenceNumberDesc(10L)).thenReturn(Optional.of(priorAnswer)); + when(messageRepository.findFirstBySession_IdOrderBySequenceNumberDesc(10L)) + .thenReturn(Optional.of(priorAnswer)); - assertThatThrownBy(() -> service.submit(1L, 10L, - command("voice".getBytes(), "audio/webm", null))) - .isInstanceOfSatisfying(DomainException.class, exception -> - assertThat(exception.getErrorCode()).isEqualTo(ApiErrorCode.SESSION_INVALID_STATE)); + assertThatThrownBy(() -> service.createVoicePlaceholder(1L, 10L, null)) + .isInstanceOfSatisfying(DomainException.class, e -> + assertThat(e.getErrorCode()).isEqualTo(ApiErrorCode.SESSION_INVALID_STATE)); - verify(storage, never()).put(any(), any(), anyLong(), any()); - verify(publisher, never()).publishToAi(any(), any(), any()); + verify(messageRepository, never()).save(any()); } - private VoiceAnswerUploadCommand command(byte[] content, String contentType, String idempotencyKey) { - return new VoiceAnswerUploadCommand( - new ByteArrayInputStream(content), - content.length, - contentType, - "answer.webm", - idempotencyKey - ); + // ── attachAudioAndRequestAnalysis ───────────────────────────────────────── + + // analyze.voice 는 여기서 직접 발행하지 않는다 — 이벤트만 내고 AFTER_COMMIT 리스너가 발행한다. + @Test + void attachAudio_setsPathAndPublishesUploadedEvent() { + InterviewMessage placeholder = voicePlaceholder(200L); + + when(messageRepository.findById(200L)).thenReturn(Optional.of(placeholder)); + + MessageResult result = service.attachAudioAndRequestAnalysis( + 1L, 10L, 200L, KEY, "audio/webm"); + + assertThat(result.id()).isEqualTo(200L); + assertThat(result.audioFilePath()).isEqualTo(KEY); + assertThat(placeholder.getAudioFilePath()).isEqualTo(KEY); + + ArgumentCaptor captor = + ArgumentCaptor.forClass(VoiceAnswerUploadedEvent.class); + verify(events).publishEvent(captor.capture()); + VoiceAnswerUploadedEvent event = captor.getValue(); + assertThat(event.userId()).isEqualTo(1L); + assertThat(event.sessionId()).isEqualTo(10L); + assertThat(event.messageId()).isEqualTo(200L); + assertThat(event.audioS3Key()).isEqualTo(KEY); + assertThat(event.contentType()).isEqualTo("audio/webm"); + } + + @Test + void attachAudio_doesNotRepublishWhenAlreadyAttached() { + InterviewMessage placeholder = voicePlaceholder(200L); + placeholder.attachAudio(KEY); + + when(messageRepository.findById(200L)).thenReturn(Optional.of(placeholder)); + + MessageResult result = service.attachAudioAndRequestAnalysis( + 1L, 10L, 200L, "interview/voice/raw/10/200-retry.webm", "audio/webm"); + + assertThat(result.audioFilePath()).isEqualTo(KEY); + verify(events, never()).publishEvent(any(VoiceAnswerUploadedEvent.class)); + } + + @Test + void attachAudio_rejectsWhenMessageNotFound() { + when(messageRepository.findById(999L)).thenReturn(Optional.empty()); + + assertThatThrownBy(() -> service.attachAudioAndRequestAnalysis( + 1L, 10L, 999L, KEY, "audio/webm")) + .isInstanceOfSatisfying(DomainException.class, e -> + assertThat(e.getErrorCode()).isEqualTo(ApiErrorCode.VOICE_MESSAGE_NOT_FOUND)); + } + + // ── failVoiceUpload ─────────────────────────────────────────────────────── + + // 업로드 실패 보상: 메시지를 지우지 않고 FAILED 로 두고 SSE 로 알린다(턴 잠김 방지). + @Test + void failVoiceUpload_marksFailedAndNotifies() { + InterviewMessage placeholder = voicePlaceholder(200L); + + when(messageRepository.findById(200L)).thenReturn(Optional.of(placeholder)); + + service.failVoiceUpload(10L, 200L); + + assertThat(placeholder.getStatus()).isEqualTo(MessageStatus.FAILED); + verify(events, org.mockito.Mockito.times(2)).publishEvent(any(RealtimeNotifyEvent.class)); + } + + @Test + void failVoiceUpload_isNoopWhenMessageMissing() { + when(messageRepository.findById(200L)).thenReturn(Optional.empty()); + + service.failVoiceUpload(10L, 200L); + + verify(events, never()).publishEvent(any()); + } + + // ── fixtures ────────────────────────────────────────────────────────────── + + private InterviewMessage voicePlaceholder(Long id) { + InterviewSession session = sessionInProgress(10L); + InterviewMessage question = InterviewMessage.interviewer(session, 1, "Tell me about ACID."); + ReflectionTestUtils.setField(question, "id", 100L); + InterviewMessage placeholder = InterviewMessage.voiceInterviewee(session, 2, question, null); + ReflectionTestUtils.setField(placeholder, "id", id); + return placeholder; } private InterviewSession sessionInProgress(Long id) { @@ -196,51 +236,4 @@ private InterviewSession sessionFixture(Long id) { ReflectionTestUtils.setField(s, "id", id); return s; } - - private RabbitMqProperties rabbitMqProperties() { - return new RabbitMqProperties( - "core", - "1", - new RabbitMqProperties.Message("application/json", "UTF-8", "X-Trace-Id"), - new RabbitMqProperties.Template(true), - new RabbitMqProperties.Exchanges(true, false, - new RabbitMqProperties.Exchanges.Names("core.ai", "ai.core", "realtime")), - new RabbitMqProperties.Queues(true, - new RabbitMqProperties.Queues.Names( - "ai.analyze.resume", - "ai.analyze.repository", - "ai.analyze.cover_letter", - "ai.generate.questions", - "ai.generate.followup", - "ai.generate.feedback", - "ai.analyze.voice", - "ai.generate.tts", - "core.callback.analysis", - "core.callback.questions", - "core.callback.feedback", - "core.callback.voice", - "core.callback.tts" - )), - new RabbitMqProperties.RoutingKeyProperties( - "analyze.resume", - "analyze.repository", - "analyze.cover_letter", - "generate.questions", - "generate.followup", - "generate.feedback", - "analyze.voice", - "generate.tts", - "callback.analysis", - "callback.questions", - "callback.feedback", - "callback.voice", - "callback.tts", - "session.notify", - "realtime.user.notify", - "realtime.document.notify" - ), - new RabbitMqProperties.DeadLetter("dlx", "dlq."), - new RabbitMqProperties.Retry(3, Duration.ofSeconds(1), 2.0, Duration.ofSeconds(10)) - ); - } } diff --git a/docs/messaging.md b/docs/messaging.md index 14aacbfd..a0960536 100644 --- a/docs/messaging.md +++ b/docs/messaging.md @@ -494,6 +494,37 @@ AI followup consumer 가 토큰 스트림 중 문장 경계마다 그 문장만 - 세그먼트 S3 키 규칙(AI·Core 공유): `interview/tts/{sessionId}/{messageId}/seg-{seq}.{ext}`. Core 는 DB 미기록, `GET …/messages/{mid}/audio/segments/{seq}?ext=` 프록시에서 소유권 검증 후 규칙으로 키 재구성(ext 화이트리스트 wav|mp3|ogg|m4a). - 상세 SSE 스펙: [`event-stream.md §3.2-1`](./event-stream.md). +### 5.14 발행 시점 규약 — 반드시 커밋 후에 발행한다 + +**Core 의 모든 작업 요청 발행(`stackup.core-to-ai`)은 DB 커밋 이후에 일어나야 한다.** +트랜잭션 안에서 발행하면 이후 커밋이 실패했을 때 AI 는 존재하지 않는 행을 대상으로 작업하고, +결과 콜백은 "not found" 로 드롭되어 **사용자 입력이 조용히 사라진다**. + +구현 패턴 — 도메인 이벤트 + `AFTER_COMMIT` 리스너: + +```java +// 1) 트랜잭션 안: DB 쓰기 + 도메인 이벤트만 +events.publishEvent(new VoiceAnswerUploadedEvent(userId, sessionId, messageId, key, contentType)); + +// 2) 커밋 후: envelope 발행 +@Transactional(readOnly = true, propagation = Propagation.REQUIRES_NEW) +@TransactionalEventListener(phase = TransactionPhase.AFTER_COMMIT) +public void onVoiceAnswerUploaded(VoiceAnswerUploadedEvent event) { … publisher.publishToAi(…); } +``` + +| Routing Key | 발행 주체 (`AFTER_COMMIT`) | 트리거 이벤트 | +|---|---|---| +| `generate.questions` | `SessionQuestionsRequester` | `SessionCreatedEvent` · `SelfIntroAnsweredEvent` | +| `generate.followup` | `SessionFollowupRequester` | `AnswerSubmittedEvent` | +| `generate.tts` | `SessionTtsRequester` | `QuestionPersistedEvent` | +| `generate.feedback` | `SessionFeedbackRequester` | `SessionEndedEvent` | +| `analyze.voice` | `VoiceAnalysisRequester` | `VoiceAnswerUploadedEvent` | +| `analyze.resume` · `analyze.repository` · `analyze.cover_letter` | `AnalysisRequestService` | `*AnalysisRequestedEvent` | + +**S3 PUT 도 트랜잭션 밖에서 한다.** 업로드(최대 25MB)가 끝날 때까지 DB 커넥션을 붙잡으면 +커넥션 풀을 잠식한다. 업로드 후 별도 트랜잭션에서 키를 붙이고, 업로드가 실패하면 이미 커밋된 +placeholder 를 `FAILED` 로 확정해 클라이언트의 턴이 잠기지 않게 보상한다. + --- ## 6. 재시도·DLQ 정책