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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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, "접근 권한이 없습니다."),
Expand Down
Original file line number Diff line number Diff line change
@@ -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());
}
}
Original file line number Diff line number Diff line change
@@ -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<String> 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);
}
}
Original file line number Diff line number Diff line change
@@ -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<String> 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)
Expand Down Expand Up @@ -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);
}
}
Original file line number Diff line number Diff line change
@@ -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
) {
}
Original file line number Diff line number Diff line change
@@ -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;
Expand Down Expand Up @@ -30,7 +30,7 @@
@RequiredArgsConstructor
public class VoiceAnswerController {

private final VoiceAnswerUploadService uploadService;
private final VoiceAnswerSubmitService submitService;
private final VoiceStreamService streamService;

@Operation(
Expand Down Expand Up @@ -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);
}

Expand Down
Loading
Loading