From e99c23d3e8f6c20694a2dabf3c6b310dc07a9396 Mon Sep 17 00:00:00 2001 From: Jaeho Date: Sun, 23 Aug 2026 03:25:53 +0900 Subject: [PATCH 1/3] =?UTF-8?q?refactor(ai):=20=EC=8B=A4=ED=8C=A8=20?= =?UTF-8?q?=EC=8B=A0=ED=98=B8=20=EA=B0=80=EB=93=9C=20=EC=A0=84=20=EC=BB=A8?= =?UTF-8?q?=EC=8A=88=EB=A8=B8=20=ED=99=95=EC=9E=A5=20+=20=ED=94=BC?= =?UTF-8?q?=EB=93=9C=EB=B0=B1=20attemptId=20=EC=97=90=EC=BD=94?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 분석 4종(resume/web/repository/cover_letter)을 공용 가드로 전환 — 도메인 분류는 공용 analysis_failed_payload 팩토리가 보존하고, expected_errors 로 일상 실패는 warning 레벨(trace_id·error_code 포함)로 기록한다. voice/tts 는 직접-발행 모델을 유지하되 멱등 마킹 이후 전 구간을 unmark_on_error 로 감싸 콜백 0건 DLQ 재주입 삼킴을 막는다(재주입은 재과금·세그먼트 재전송 동반 — 수동 복구 전용). feedback 은 요청의 attemptId 를 성공/FAILED 콜백에 에코 — Core 가 대체된 이전 시도의 지연 FAILED 를 드롭하는 근거. 테스트 +6. --- ai/CLAUDE.md | 8 +- ai/src/ai_server/CLAUDE.md | 10 +- .../consumers/cover_letter_consumer.py | 144 +++++-------- .../messaging/consumers/failure_signal.py | 88 ++++++-- .../messaging/consumers/feedback_consumer.py | 2 + .../consumers/repository_consumer.py | 148 +++++--------- .../messaging/consumers/resume_consumer.py | 140 +++++-------- .../messaging/consumers/tts_consumer.py | 85 ++++---- .../messaging/consumers/voice_consumer.py | 191 +++++++++--------- .../messaging/consumers/web_consumer.py | 132 ++++-------- ai/src/ai_server/model/messages/feedback.py | 5 + ai/tests/test_cover_letter_consumer.py | 18 ++ ai/tests/test_feedback_consumer.py | 3 + ai/tests/test_resume_consumer.py | 24 +++ ai/tests/test_tts_consumer.py | 54 +++++ ai/tests/test_voice_consumer.py | 23 +++ ai/tests/test_web_consumer.py | 18 ++ 17 files changed, 565 insertions(+), 528 deletions(-) diff --git a/ai/CLAUDE.md b/ai/CLAUDE.md index 452e9e77..b710103a 100644 --- a/ai/CLAUDE.md +++ b/ai/CLAUDE.md @@ -397,7 +397,13 @@ docker run --env-file .env -p 8000:8000 stackup-ai payload 반환)와 `_failed_payload(req, exc)` 팩토리만 구현한다. `errorCode` 분류는 consumer 별 계약 유지 — questions/followup: `GENERATION_SCHEMA_INVALID`(TypeError, retriable=false) | `GENERATION_FAILED`, feedback: 동일 스키마 코드 | `UNEXPECTED`. `errorMessage` 는 공용 - `format_error_message`(`ExcType: msg`, 500자 상한). Core 쪽 처리는 + `format_error_message`(`ExcType: msg`, 500자 상한). **F5/F6 확장**: 분석 4개 + (resume/web/repository/cover_letter)도 같은 가드로 전환(도메인 에러의 코드·retriable 분류는 + 각 `_failed_payload` 팩토리가 보존, `action="analyze"` 로 기존 로그 키 유지), voice/tts 는 + 직접-발행 모델을 유지하되 마킹 이후 전 구간을 `unmark_on_error` 로 감쌌다(재주입은 + 합성 재과금·세그먼트 재전송을 동반 — 수동 복구 전용). feedback 은 요청의 + `attemptId` 를 성공/FAILED 콜백에 에코 — Core 가 대체된 이전 시도의 지연 FAILED 를 드롭하는 + 근거(messaging.md §5.10/§5.11). Core 쪽 처리는 [`backend/CLAUDE.md`](../backend/CLAUDE.md) 참고. (질문 풀 RAG `questions_rag_timeout_sec` 1.5s 하드 타임아웃과 `llm_pro_timeout_sec` 30s/`llm_flash_timeout_sec` 10s 요청 타임아웃은 이전 리팩터에서 도입되어 유지된다.) diff --git a/ai/src/ai_server/CLAUDE.md b/ai/src/ai_server/CLAUDE.md index 0013f983..6bed8ce3 100644 --- a/ai/src/ai_server/CLAUDE.md +++ b/ai/src/ai_server/CLAUDE.md @@ -76,10 +76,12 @@ class QuestionPoolCallbackPayload(BaseModel): 콜백 발행은 `publisher.py`, 멱등은 `idempotency.py`(`LruIdempotencyStore`), RealTime 직접 발행은 `progress.py`(분석 진행, user 채널)·`session_notify.py`(델타/오디오/질문 풀·피드백 생성 진행, 세션 채널) - 모든 consumer는 envelope parsing → trace_context → 비즈니스 핸들러 호출 패턴 -- 생성 계열 3개(questions/followup/feedback)는 공용 가드 `consumers/failure_signal.py: - consume_with_failure_signal` 경유 — consumer 는 `_process(envelope)`(성공 payload 반환)와 - `_failed_payload(req, exc)` 팩토리만 구현하고, 파싱→멱등→전 구간 가드→FAILED 콜백/성공 발행→ - unmark 는 가드가 책임진다 ([`/docs/messaging.md §6`](../../../docs/messaging.md) AI Server 절) +- 생성 3개(questions/followup/feedback) + 분석 4개(resume/web/repository/cover_letter)는 공용 가드 + `consumers/failure_signal.py: consume_with_failure_signal` 경유 — consumer 는 + `_process(envelope)`(성공 payload 반환)와 `_failed_payload(req, exc)` 팩토리만 구현하고, + 파싱→멱등→전 구간 가드→FAILED 콜백/성공 발행→unmark 는 가드가 책임진다. voice/tts 는 + 직접-발행 모델 유지 + 마킹 이후 전 구간을 `unmark_on_error` 로 래핑 + ([`/docs/messaging.md §6`](../../../docs/messaging.md) AI Server 절) ```python # messaging/consumers/resume_consumer.py (패턴) diff --git a/ai/src/ai_server/messaging/consumers/cover_letter_consumer.py b/ai/src/ai_server/messaging/consumers/cover_letter_consumer.py index 827bfde3..55f22014 100644 --- a/ai/src/ai_server/messaging/consumers/cover_letter_consumer.py +++ b/ai/src/ai_server/messaging/consumers/cover_letter_consumer.py @@ -4,6 +4,11 @@ from aio_pika.abc import AbstractIncomingMessage from ai_server.analyzer.resume_analyzer import ResumeAnalyzeError, ResumeAnalyzer +from ai_server.messaging.consumers.failure_signal import ( + analysis_done_fields, + analysis_failed_payload, + consume_with_failure_signal, +) from ai_server.messaging.idempotency import LruIdempotencyStore from ai_server.messaging.progress import AnalysisProgressNotifier from ai_server.messaging.publisher import CallbackPublisher @@ -35,112 +40,49 @@ def __init__( self._progress = progress_notifier async def handle(self, message: AbstractIncomingMessage) -> None: - async with message.process(requeue=False): - try: - envelope = Envelope[CoverLetterAnalyzeRequest].model_validate_json( - message.body - ) - except Exception as exc: # parse error → DLQ-ready (auto NACK on raise) - log.error( - "cover_letter.parse.failed", - error=str(exc), - delivery_tag=message.delivery_tag, - ) - raise - - if self._idempotency.is_seen_then_mark(envelope.message_id): - log.info( - "cover_letter.idempotent.skip", - message_id=envelope.message_id, - trace_id=envelope.trace_id, - ) - return - - req = envelope.payload - log.info( - "cover_letter.analyze.start", - message_id=envelope.message_id, - cover_letter_id=req.cover_letter_id, - trace_id=envelope.trace_id, - ) - - payload = await self._run_and_build_payload( - req, envelope.trace_id, user_id=envelope.context.user_id - ) - - await self._publisher.publish( - routing_key=self._callback_routing_key, - message_type="callback.analysis", - payload=payload, - trace_id=envelope.trace_id, - correlation_id=envelope.message_id, - context=envelope.context, - ) - log.info( - "cover_letter.analyze.done", - message_id=envelope.message_id, - cover_letter_id=req.cover_letter_id, - status=payload.status, - trace_id=envelope.trace_id, - ) + await consume_with_failure_signal( + message, + domain="cover_letter", + action="analyze", + envelope_type=Envelope[CoverLetterAnalyzeRequest], + idempotency=self._idempotency, + publisher=self._publisher, + routing_key=self._callback_routing_key, + message_type="callback.analysis", + process=self._process, + failed_payload=self._failed_payload, + done_fields=analysis_done_fields, + expected_errors=(ResumeAnalyzeError,), + ) - async def _run_and_build_payload( - self, - req: CoverLetterAnalyzeRequest, - trace_id: str, - *, - user_id: int | None, + async def _process( + self, envelope: Envelope[CoverLetterAnalyzeRequest] ) -> AnalysisCallbackPayload: + req = envelope.payload + log.info( + "cover_letter.analyze.start", + message_id=envelope.message_id, + cover_letter_id=req.cover_letter_id, + trace_id=envelope.trace_id, + ) progress = ( self._progress.emitter_for( - user_id=user_id, + user_id=envelope.context.user_id, target_type="COVER_LETTER", target_id=req.cover_letter_id, - trace_id=trace_id, + trace_id=envelope.trace_id, ) if self._progress is not None else None ) - try: - # ResumeAnalyzer 의 resume_id/file_path 는 각각 식별자/추출 locator 로 일반화돼 있어 - # 자소서는 cover_letter_id 와 inline content 를 그대로 넘긴다. - result = await self._analyzer.analyze( - resume_id=req.cover_letter_id, - file_path=req.content, - analyzed_document_id=req.analyzed_document_id, - progress=progress, - ) - except ResumeAnalyzeError as err: - log.warning( - "cover_letter.analyze.domain_failed", - cover_letter_id=req.cover_letter_id, - code=err.code, - retriable=err.retriable, - trace_id=trace_id, - ) - return AnalysisCallbackPayload( - target_type="COVER_LETTER", - target_id=req.cover_letter_id, - status="FAILED", - error_code=err.code, - error_message=err.message, - retriable=err.retriable, - ) - except Exception as exc: - log.exception( - "cover_letter.analyze.unexpected_failed", - cover_letter_id=req.cover_letter_id, - trace_id=trace_id, - ) - return AnalysisCallbackPayload( - target_type="COVER_LETTER", - target_id=req.cover_letter_id, - status="FAILED", - error_code="UNEXPECTED", - error_message=str(exc), - retriable=True, - ) - + # ResumeAnalyzer 의 resume_id/file_path 는 각각 식별자/추출 locator 로 일반화돼 있어 + # 자소서는 cover_letter_id 와 inline content 를 그대로 넘긴다. + result = await self._analyzer.analyze( + resume_id=req.cover_letter_id, + file_path=req.content, + analyzed_document_id=req.analyzed_document_id, + progress=progress, + ) return AnalysisCallbackPayload( target_type="COVER_LETTER", target_id=req.cover_letter_id, @@ -150,3 +92,13 @@ async def _run_and_build_payload( document_path=result.document_path, embedding_chunk_count=result.embedding_chunk_count, ) + + def _failed_payload( + self, req: CoverLetterAnalyzeRequest, exc: Exception + ) -> AnalysisCallbackPayload: + return analysis_failed_payload( + target_type="COVER_LETTER", + target_id=req.cover_letter_id, + exc=exc, + domain_error=ResumeAnalyzeError, + ) diff --git a/ai/src/ai_server/messaging/consumers/failure_signal.py b/ai/src/ai_server/messaging/consumers/failure_signal.py index 95fa585b..acf0fe3c 100644 --- a/ai/src/ai_server/messaging/consumers/failure_signal.py +++ b/ai/src/ai_server/messaging/consumers/failure_signal.py @@ -1,6 +1,7 @@ from __future__ import annotations -from collections.abc import Awaitable, Callable +import contextlib +from collections.abc import AsyncIterator, Awaitable, Callable from typing import Any, TypeVar import structlog @@ -10,6 +11,7 @@ from ai_server.messaging.idempotency import LruIdempotencyStore from ai_server.messaging.publisher import CallbackPublisher from ai_server.model.envelope import Envelope +from ai_server.model.messages.analyze import AnalysisCallbackPayload log = structlog.get_logger(__name__) @@ -30,6 +32,54 @@ def classify_failure(exc: Exception, *, unexpected_code: str) -> tuple[str, bool return unexpected_code, True +@contextlib.asynccontextmanager +async def unmark_on_error( + idempotency: LruIdempotencyStore, message_id: str +) -> AsyncIterator[None]: + """멱등 마킹 이후 구간의 어떤 예외든 unmark 후 전파 — 발행 실패뿐 아니라 그 앞의 + 합성·계산 단계에서 죽어도, 콜백 0건으로 DLQ 에 간 메시지의 재주입이 duplicate skip 으로 + 삼켜지지 않게 한다. 전면 가드 전환이 부담스러운 직접-발행 모델 컨슈머(voice/tts)용. + 주의: 재주입은 pre-publish 부수효과(합성·STT 재과금, 라이브 세그먼트 재전송)를 재실행한다 — + 운영자 수동 복구 수단으로만 쓴다 (docs/messaging.md §6).""" + try: + yield + except Exception: + idempotency.unmark(message_id) + raise + + +def analysis_failed_payload( + *, + target_type: str, + target_id: int, + exc: Exception, + domain_error: type[Exception], +) -> AnalysisCallbackPayload: + """분석 4종(resume/web/repository/cover_letter) 공용 FAILED payload — + 도메인 에러는 코드·retriable 분류를 보존하고, 그 외는 UNEXPECTED/retriable=true.""" + if isinstance(exc, domain_error): + return AnalysisCallbackPayload( + target_type=target_type, + target_id=target_id, + status="FAILED", + error_code=exc.code, # type: ignore[attr-defined] + error_message=exc.message, # type: ignore[attr-defined] + retriable=exc.retriable, # type: ignore[attr-defined] + ) + return AnalysisCallbackPayload( + target_type=target_type, + target_id=target_id, + status="FAILED", + error_code="UNEXPECTED", + error_message=format_error_message(exc), + retriable=True, + ) + + +def analysis_done_fields(payload: Any) -> dict[str, Any]: + return {"target_id": payload.target_id, "status": payload.status} + + async def consume_with_failure_signal( message: AbstractIncomingMessage, *, @@ -42,6 +92,8 @@ async def consume_with_failure_signal( process: Callable[[Envelope[ReqT]], Awaitable[BaseModel]], failed_payload: Callable[[ReqT, Exception], BaseModel], done_fields: Callable[[Any], dict[str, Any]] | None = None, + action: str = "generate", + expected_errors: tuple[type[Exception], ...] = (), ) -> None: """생성 컨슈머(질문 풀·꼬리질문·피드백) 공용 실패 신호 가드. @@ -84,17 +136,29 @@ async def _publish(payload: BaseModel) -> None: context=envelope.context, ) + log_ids: dict[str, Any] = { + "message_id": envelope.message_id, + "trace_id": envelope.trace_id, + } session_id = getattr(envelope.payload, "session_id", None) + if session_id is not None: + log_ids["session_id"] = session_id try: payload = await process(envelope) except Exception as exc: # noqa: BLE001 - log.exception( - f"{domain}.generate.failed", - message_id=envelope.message_id, - session_id=session_id, - trace_id=envelope.trace_id, - ) + if expected_errors and isinstance(exc, expected_errors): + # 예상된 도메인 실패(빈 PDF·404 URL 등 일상 입력)는 traceback 없는 warning — + # ERROR 레벨 스택트레이스로 일상 실패가 알람을 울리지 않게 한다. + log.warning( + f"{domain}.{action}.failed", + error_code=getattr(exc, "code", None), + retriable=getattr(exc, "retriable", None), + error=str(exc), + **log_ids, + ) + else: + log.exception(f"{domain}.{action}.failed", **log_ids) try: # 팩토리 자체가 죽어도(검증 오류 등) 같은 안전망을 태운다 — # try 밖이면 unmark 없이 DLQ 로 가서 재주입이 duplicate skip 으로 삼켜진다. @@ -103,9 +167,7 @@ async def _publish(payload: BaseModel) -> None: except Exception: # noqa: BLE001 log.exception( f"{domain}.failed_callback.publish_failed", - message_id=envelope.message_id, - session_id=session_id, - trace_id=envelope.trace_id, + **log_ids, ) idempotency.unmark(envelope.message_id) raise exc @@ -117,9 +179,7 @@ async def _publish(payload: BaseModel) -> None: idempotency.unmark(envelope.message_id) raise log.info( - f"{domain}.generate.done", - message_id=envelope.message_id, - session_id=session_id, - trace_id=envelope.trace_id, + f"{domain}.{action}.done", + **log_ids, **(done_fields(payload) if done_fields is not None else {}), ) diff --git a/ai/src/ai_server/messaging/consumers/feedback_consumer.py b/ai/src/ai_server/messaging/consumers/feedback_consumer.py index 0a01dc15..60b96938 100644 --- a/ai/src/ai_server/messaging/consumers/feedback_consumer.py +++ b/ai/src/ai_server/messaging/consumers/feedback_consumer.py @@ -228,6 +228,7 @@ async def _tracked(coro: Awaitable[T]) -> T: panel_breakdown=result.panel_breakdown, answer_coaching=answer_coaching, report_s3_key=None, + attempt_id=req.attempt_id, ) return payload @@ -245,6 +246,7 @@ def _failed_payload( error_code=error_code, error_message=format_error_message(exc), retriable=retriable, + attempt_id=req.attempt_id, ) async def _emit_progress( diff --git a/ai/src/ai_server/messaging/consumers/repository_consumer.py b/ai/src/ai_server/messaging/consumers/repository_consumer.py index 60ef7cc0..235d6d3e 100644 --- a/ai/src/ai_server/messaging/consumers/repository_consumer.py +++ b/ai/src/ai_server/messaging/consumers/repository_consumer.py @@ -7,6 +7,11 @@ RepositoryAnalyzeError, RepositoryAnalyzer, ) +from ai_server.messaging.consumers.failure_signal import ( + analysis_done_fields, + analysis_failed_payload, + consume_with_failure_signal, +) from ai_server.messaging.idempotency import LruIdempotencyStore from ai_server.messaging.progress import AnalysisProgressNotifier from ai_server.messaging.publisher import CallbackPublisher @@ -36,115 +41,52 @@ def __init__( self._progress = progress_notifier async def handle(self, message: AbstractIncomingMessage) -> None: - async with message.process(requeue=False): - try: - envelope = Envelope[RepositoryAnalyzeRequest].model_validate_json( - message.body - ) - except Exception as exc: - log.error( - "repository.parse.failed", - error=str(exc), - delivery_tag=message.delivery_tag, - ) - raise - - if self._idempotency.is_seen_then_mark(envelope.message_id): - log.info( - "repository.idempotent.skip", - message_id=envelope.message_id, - trace_id=envelope.trace_id, - ) - return - - req = envelope.payload - user_id = envelope.context.user_id - log.info( - "repository.analyze.start", - message_id=envelope.message_id, - repository_id=req.repository_id, - repo_full_name=req.repo_full_name, - user_id=user_id, - trace_id=envelope.trace_id, - ) - - payload = await self._run_and_build_payload( - req, user_id=user_id, trace_id=envelope.trace_id - ) - - await self._publisher.publish( - routing_key=self._callback_routing_key, - message_type="callback.analysis", - payload=payload, - trace_id=envelope.trace_id, - correlation_id=envelope.message_id, - context=envelope.context, - ) - log.info( - "repository.analyze.done", - message_id=envelope.message_id, - repository_id=req.repository_id, - status=payload.status, - trace_id=envelope.trace_id, - ) + await consume_with_failure_signal( + message, + domain="repository", + action="analyze", + envelope_type=Envelope[RepositoryAnalyzeRequest], + idempotency=self._idempotency, + publisher=self._publisher, + routing_key=self._callback_routing_key, + message_type="callback.analysis", + process=self._process, + failed_payload=self._failed_payload, + done_fields=analysis_done_fields, + expected_errors=(RepositoryAnalyzeError,), + ) - async def _run_and_build_payload( - self, - req: RepositoryAnalyzeRequest, - *, - user_id: int | None, - trace_id: str, + async def _process( + self, envelope: Envelope[RepositoryAnalyzeRequest] ) -> AnalysisCallbackPayload: + req = envelope.payload + user_id = envelope.context.user_id + log.info( + "repository.analyze.start", + message_id=envelope.message_id, + repository_id=req.repository_id, + repo_full_name=req.repo_full_name, + user_id=user_id, + trace_id=envelope.trace_id, + ) progress = ( self._progress.emitter_for( user_id=user_id, target_type="REPOSITORY", target_id=req.repository_id, - trace_id=trace_id, + trace_id=envelope.trace_id, ) if self._progress is not None else None ) - try: - result = await self._analyzer.analyze( - repository_id=req.repository_id, - repo_full_name=req.repo_full_name, - default_branch=req.default_branch, - user_id=user_id, - analyzed_document_id=req.analyzed_document_id, - progress=progress, - ) - except RepositoryAnalyzeError as err: - log.warning( - "repository.analyze.domain_failed", - repository_id=req.repository_id, - code=err.code, - retriable=err.retriable, - trace_id=trace_id, - ) - return AnalysisCallbackPayload( - target_type="REPOSITORY", - target_id=req.repository_id, - status="FAILED", - error_code=err.code, - error_message=err.message, - retriable=err.retriable, - ) - except Exception as exc: - log.exception( - "repository.analyze.unexpected_failed", - repository_id=req.repository_id, - trace_id=trace_id, - ) - return AnalysisCallbackPayload( - target_type="REPOSITORY", - target_id=req.repository_id, - status="FAILED", - error_code="UNEXPECTED", - error_message=str(exc), - retriable=True, - ) - + result = await self._analyzer.analyze( + repository_id=req.repository_id, + repo_full_name=req.repo_full_name, + default_branch=req.default_branch, + user_id=user_id, + analyzed_document_id=req.analyzed_document_id, + progress=progress, + ) return AnalysisCallbackPayload( target_type="REPOSITORY", target_id=req.repository_id, @@ -154,3 +96,13 @@ async def _run_and_build_payload( document_path=result.document_path, embedding_chunk_count=result.embedding_chunk_count, ) + + def _failed_payload( + self, req: RepositoryAnalyzeRequest, exc: Exception + ) -> AnalysisCallbackPayload: + return analysis_failed_payload( + target_type="REPOSITORY", + target_id=req.repository_id, + exc=exc, + domain_error=RepositoryAnalyzeError, + ) diff --git a/ai/src/ai_server/messaging/consumers/resume_consumer.py b/ai/src/ai_server/messaging/consumers/resume_consumer.py index f22ed3a0..50ad3ac2 100644 --- a/ai/src/ai_server/messaging/consumers/resume_consumer.py +++ b/ai/src/ai_server/messaging/consumers/resume_consumer.py @@ -4,6 +4,11 @@ from aio_pika.abc import AbstractIncomingMessage from ai_server.analyzer.resume_analyzer import ResumeAnalyzeError, ResumeAnalyzer +from ai_server.messaging.consumers.failure_signal import ( + analysis_done_fields, + analysis_failed_payload, + consume_with_failure_signal, +) from ai_server.messaging.idempotency import LruIdempotencyStore from ai_server.messaging.progress import AnalysisProgressNotifier from ai_server.messaging.publisher import CallbackPublisher @@ -33,110 +38,47 @@ def __init__( self._progress = progress_notifier async def handle(self, message: AbstractIncomingMessage) -> None: - async with message.process(requeue=False): - try: - envelope = Envelope[ResumeAnalyzeRequest].model_validate_json( - message.body - ) - except Exception as exc: # parse error → DLQ-ready (auto NACK on raise) - log.error( - "resume.parse.failed", - error=str(exc), - delivery_tag=message.delivery_tag, - ) - raise - - if self._idempotency.is_seen_then_mark(envelope.message_id): - log.info( - "resume.idempotent.skip", - message_id=envelope.message_id, - trace_id=envelope.trace_id, - ) - return - - req = envelope.payload - log.info( - "resume.analyze.start", - message_id=envelope.message_id, - resume_id=req.resume_id, - trace_id=envelope.trace_id, - ) - - payload = await self._run_and_build_payload( - req, envelope.trace_id, user_id=envelope.context.user_id - ) - - await self._publisher.publish( - routing_key=self._callback_routing_key, - message_type="callback.analysis", - payload=payload, - trace_id=envelope.trace_id, - correlation_id=envelope.message_id, - context=envelope.context, - ) - log.info( - "resume.analyze.done", - message_id=envelope.message_id, - resume_id=req.resume_id, - status=payload.status, - trace_id=envelope.trace_id, - ) + await consume_with_failure_signal( + message, + domain="resume", + action="analyze", + envelope_type=Envelope[ResumeAnalyzeRequest], + idempotency=self._idempotency, + publisher=self._publisher, + routing_key=self._callback_routing_key, + message_type="callback.analysis", + process=self._process, + failed_payload=self._failed_payload, + done_fields=analysis_done_fields, + expected_errors=(ResumeAnalyzeError,), + ) - async def _run_and_build_payload( - self, - req: ResumeAnalyzeRequest, - trace_id: str, - *, - user_id: int | None, + async def _process( + self, envelope: Envelope[ResumeAnalyzeRequest] ) -> AnalysisCallbackPayload: + req = envelope.payload + log.info( + "resume.analyze.start", + message_id=envelope.message_id, + resume_id=req.resume_id, + trace_id=envelope.trace_id, + ) progress = ( self._progress.emitter_for( - user_id=user_id, + user_id=envelope.context.user_id, target_type="RESUME", target_id=req.resume_id, - trace_id=trace_id, + trace_id=envelope.trace_id, ) if self._progress is not None else None ) - try: - result = await self._analyzer.analyze( - resume_id=req.resume_id, - file_path=req.file_path, - analyzed_document_id=req.analyzed_document_id, - progress=progress, - ) - except ResumeAnalyzeError as err: - log.warning( - "resume.analyze.domain_failed", - resume_id=req.resume_id, - code=err.code, - retriable=err.retriable, - trace_id=trace_id, - ) - return AnalysisCallbackPayload( - target_type="RESUME", - target_id=req.resume_id, - status="FAILED", - error_code=err.code, - error_message=err.message, - retriable=err.retriable, - ) - except Exception as exc: - log.exception( - "resume.analyze.unexpected_failed", - resume_id=req.resume_id, - trace_id=trace_id, - ) - return AnalysisCallbackPayload( - target_type="RESUME", - target_id=req.resume_id, - status="FAILED", - error_code="UNEXPECTED", - error_message=str(exc), - retriable=True, - ) - + result = await self._analyzer.analyze( + resume_id=req.resume_id, + file_path=req.file_path, + analyzed_document_id=req.analyzed_document_id, + progress=progress, + ) return AnalysisCallbackPayload( target_type="RESUME", target_id=req.resume_id, @@ -146,3 +88,13 @@ async def _run_and_build_payload( document_path=result.document_path, embedding_chunk_count=result.embedding_chunk_count, ) + + def _failed_payload( + self, req: ResumeAnalyzeRequest, exc: Exception + ) -> AnalysisCallbackPayload: + return analysis_failed_payload( + target_type="RESUME", + target_id=req.resume_id, + exc=exc, + domain_error=ResumeAnalyzeError, + ) diff --git a/ai/src/ai_server/messaging/consumers/tts_consumer.py b/ai/src/ai_server/messaging/consumers/tts_consumer.py index c76dceab..a56ca617 100644 --- a/ai/src/ai_server/messaging/consumers/tts_consumer.py +++ b/ai/src/ai_server/messaging/consumers/tts_consumer.py @@ -7,6 +7,7 @@ from aio_pika.abc import AbstractIncomingMessage from ai_server.chain.sentence_split import next_sentences +from ai_server.messaging.consumers.failure_signal import unmark_on_error from ai_server.messaging.idempotency import LruIdempotencyStore from ai_server.messaging.publisher import CallbackPublisher from ai_server.messaging.session_notify import SessionRealtimeNotifier @@ -84,57 +85,61 @@ async def handle(self, message: AbstractIncomingMessage) -> None: log.info("tts.idempotent.skip", message_id=envelope.message_id) return - req = envelope.payload - log.info( - "tts.synthesize.start", - message_id=envelope.message_id, - session_id=req.session_id, - interview_message_id=req.message_id, - streaming=self._notifier is not None, - trace_id=envelope.trace_id, - ) + # 마킹 이후 어떤 예외든 unmark — 콜백 0건 DLQ 재주입 삼킴 방지 (F6). + async with unmark_on_error(self._idempotency, envelope.message_id): + req = envelope.payload + log.info( + "tts.synthesize.start", + message_id=envelope.message_id, + session_id=req.session_id, + interview_message_id=req.message_id, + streaming=self._notifier is not None, + trace_id=envelope.trace_id, + ) - if self._notifier is not None: - result = await self._synthesize_segmented(envelope, req) - else: - result = await self._synthesize_once(envelope, req) - if result is None: - return # 실패 — 콜백은 내부에서 이미 발행 + if self._notifier is not None: + result = await self._synthesize_segmented(envelope, req) + else: + result = await self._synthesize_once(envelope, req) + if result is None: + return # 실패 — 콜백은 내부에서 이미 발행 - key = self._build_key(req.session_id, req.message_id, result.content_type) - try: - await self._storage.put_bytes( - key, result.audio_bytes, content_type=result.content_type + key = self._build_key( + req.session_id, req.message_id, result.content_type ) - except Exception as exc: - log.error("tts.storage.failed", error=str(exc), key=key) + try: + await self._storage.put_bytes( + key, result.audio_bytes, content_type=result.content_type + ) + except Exception as exc: + log.error("tts.storage.failed", error=str(exc), key=key) + await self._publish( + envelope, + TtsCallbackPayload( + session_id=req.session_id, + message_id=req.message_id, + status="FAILED", + error_code="TTS_STORAGE_FAILED", + ), + ) + return + await self._publish( envelope, TtsCallbackPayload( session_id=req.session_id, message_id=req.message_id, - status="FAILED", - error_code="TTS_STORAGE_FAILED", + status="SUCCEEDED", + audio_key=key, + duration_sec=result.duration_sec, ), ) - return - - await self._publish( - envelope, - TtsCallbackPayload( + log.info( + "tts.synthesize.done", session_id=req.session_id, - message_id=req.message_id, - status="SUCCEEDED", - audio_key=key, - duration_sec=result.duration_sec, - ), - ) - log.info( - "tts.synthesize.done", - session_id=req.session_id, - interview_message_id=req.message_id, - key=key, - ) + interview_message_id=req.message_id, + key=key, + ) async def _synthesize_once( self, envelope: Envelope[GenerateTtsRequest], req: GenerateTtsRequest diff --git a/ai/src/ai_server/messaging/consumers/voice_consumer.py b/ai/src/ai_server/messaging/consumers/voice_consumer.py index d90408ff..0b73b15c 100644 --- a/ai/src/ai_server/messaging/consumers/voice_consumer.py +++ b/ai/src/ai_server/messaging/consumers/voice_consumer.py @@ -7,6 +7,7 @@ from aio_pika.abc import AbstractIncomingMessage from ai_server.core.client import CoreClient +from ai_server.messaging.consumers.failure_signal import unmark_on_error from ai_server.messaging.idempotency import LruIdempotencyStore from ai_server.messaging.publisher import CallbackPublisher from ai_server.model.envelope import Envelope @@ -69,96 +70,107 @@ async def handle(self, message: AbstractIncomingMessage) -> None: log.info("voice.idempotent.skip", message_id=envelope.message_id) return - req = envelope.payload - log.info( - "voice.analyze.start", - message_id=envelope.message_id, - session_id=req.session_id, - interview_message_id=req.message_id, - key=req.audio_s3_key, - trace_id=envelope.trace_id, - ) - - try: - audio_bytes = await self._storage.get_bytes(req.audio_s3_key) - except Exception as exc: - log.error("voice.storage.failed", error=str(exc), key=req.audio_s3_key) - await self._publish_failed(envelope, req, code="AUDIO_FETCH_FAILED") - return - - started = None - try: - started = time.perf_counter() - result = await self._stt.transcribe( - audio_bytes=audio_bytes, - content_type=req.content_type, - hint=req.previous_question_text, - ) - self._record_stt_log( - req, - envelope, - latency_ms=_elapsed_ms(started), - status="SUCCESS", - error_message=None, - ) - except SttError as exc: - self._record_stt_log( - req, - envelope, - latency_ms=_elapsed_ms(started), - status="FAILED", - error_message=exc.message, - ) - log.error( - "voice.stt.failed", - error=str(exc), - code=exc.code, + # 마킹 이후 어떤 예외든 unmark — 콜백 0건 DLQ 재주입 삼킴 방지 (F6). + async with unmark_on_error(self._idempotency, envelope.message_id): + req = envelope.payload + log.info( + "voice.analyze.start", + message_id=envelope.message_id, session_id=req.session_id, + interview_message_id=req.message_id, + key=req.audio_s3_key, + trace_id=envelope.trace_id, ) - await self._publish_failed(envelope, req, code=exc.code) - return - except Exception as exc: - self._record_stt_log( - req, - envelope, - latency_ms=_elapsed_ms(started), - status="FAILED", - error_message=str(exc), + + try: + audio_bytes = await self._storage.get_bytes(req.audio_s3_key) + except Exception as exc: + log.error( + "voice.storage.failed", error=str(exc), key=req.audio_s3_key + ) + await self._publish_failed(envelope, req, code="AUDIO_FETCH_FAILED") + return + + started = None + try: + started = time.perf_counter() + result = await self._stt.transcribe( + audio_bytes=audio_bytes, + content_type=req.content_type, + hint=req.previous_question_text, + ) + self._record_stt_log( + req, + envelope, + latency_ms=_elapsed_ms(started), + status="SUCCESS", + error_message=None, + ) + except SttError as exc: + self._record_stt_log( + req, + envelope, + latency_ms=_elapsed_ms(started), + status="FAILED", + error_message=exc.message, + ) + log.error( + "voice.stt.failed", + error=str(exc), + code=exc.code, + session_id=req.session_id, + ) + await self._publish_failed(envelope, req, code=exc.code) + return + except Exception as exc: + self._record_stt_log( + req, + envelope, + latency_ms=_elapsed_ms(started), + status="FAILED", + error_message=str(exc), + ) + log.error( + "voice.stt.unexpected", + error=str(exc), + session_id=req.session_id, + ) + await self._publish_failed( + envelope, req, code="TRANSCRIPTION_FAILED" + ) + return + + metrics = analyze(result, filler_pattern=self._filler_pattern) + + payload = VoiceCallbackPayload( + session_id=req.session_id, + interview_message_id=req.message_id, + transcript=result.text, + speaking_rate_wpm=metrics.speaking_rate_wpm, + silence_duration_sec=metrics.silence_duration_sec, + filler_word_counts=metrics.filler_word_counts, + pronunciation_accuracy=metrics.pronunciation_accuracy, + error_code=None, ) - log.error( - "voice.stt.unexpected", error=str(exc), session_id=req.session_id + await self._publish_callback(envelope, payload) + log.info( + "voice.analyze.done", + message_id=envelope.message_id, + session_id=req.session_id, + interview_message_id=req.message_id, + wpm=metrics.speaking_rate_wpm, + trace_id=envelope.trace_id, ) - await self._publish_failed(envelope, req, code="TRANSCRIPTION_FAILED") - return - metrics = analyze(result, filler_pattern=self._filler_pattern) - - payload = VoiceCallbackPayload( - session_id=req.session_id, - interview_message_id=req.message_id, - transcript=result.text, - speaking_rate_wpm=metrics.speaking_rate_wpm, - silence_duration_sec=metrics.silence_duration_sec, - filler_word_counts=metrics.filler_word_counts, - pronunciation_accuracy=metrics.pronunciation_accuracy, - error_code=None, - ) - await self._publisher.publish( - routing_key=self._callback_routing_key, - message_type="callback.voice", - payload=payload, - trace_id=envelope.trace_id, - correlation_id=envelope.message_id, - context=envelope.context, - ) - log.info( - "voice.analyze.done", - message_id=envelope.message_id, - session_id=req.session_id, - interview_message_id=req.message_id, - wpm=metrics.speaking_rate_wpm, - trace_id=envelope.trace_id, - ) + async def _publish_callback(self, envelope, payload: VoiceCallbackPayload) -> None: + await self._publisher.publish( + routing_key=self._callback_routing_key, + message_type="callback.voice", + payload=payload, + trace_id=envelope.trace_id, + correlation_id=envelope.message_id, + context=envelope.context, + ) async def _publish_failed( self, envelope, req: AnalyzeVoiceRequest, *, code: str @@ -173,14 +185,7 @@ async def _publish_failed( pronunciation_accuracy=None, error_code=code, ) - await self._publisher.publish( - routing_key=self._callback_routing_key, - message_type="callback.voice", - payload=payload, - trace_id=envelope.trace_id, - correlation_id=envelope.message_id, - context=envelope.context, - ) + await self._publish_callback(envelope, payload) def _record_stt_log( self, diff --git a/ai/src/ai_server/messaging/consumers/web_consumer.py b/ai/src/ai_server/messaging/consumers/web_consumer.py index 830b4a89..e50ada06 100644 --- a/ai/src/ai_server/messaging/consumers/web_consumer.py +++ b/ai/src/ai_server/messaging/consumers/web_consumer.py @@ -7,6 +7,11 @@ WebResumeAnalyzeError, WebResumeAnalyzer, ) +from ai_server.messaging.consumers.failure_signal import ( + analysis_done_fields, + analysis_failed_payload, + consume_with_failure_signal, +) from ai_server.messaging.idempotency import LruIdempotencyStore from ai_server.messaging.publisher import CallbackPublisher from ai_server.model.envelope import Envelope @@ -33,96 +38,37 @@ def __init__( self._callback_routing_key = callback_routing_key async def handle(self, message: AbstractIncomingMessage) -> None: - async with message.process(requeue=False): - try: - envelope = Envelope[WebResumeAnalyzeRequest].model_validate_json( - message.body - ) - except Exception as exc: - log.error( - "web_resume.parse.failed", - error=str(exc), - delivery_tag=message.delivery_tag, - ) - raise - - if self._idempotency.is_seen_then_mark(envelope.message_id): - log.info( - "web_resume.idempotent.skip", - message_id=envelope.message_id, - trace_id=envelope.trace_id, - ) - return - - req = envelope.payload - log.info( - "web_resume.analyze.start", - message_id=envelope.message_id, - resume_id=req.resume_id, - url=req.url, - trace_id=envelope.trace_id, - ) - - payload = await self._run_and_build_payload(req, envelope.trace_id) - - await self._publisher.publish( - routing_key=self._callback_routing_key, - message_type="callback.analysis", - payload=payload, - trace_id=envelope.trace_id, - correlation_id=envelope.message_id, - context=envelope.context, - ) - log.info( - "web_resume.analyze.done", - message_id=envelope.message_id, - resume_id=req.resume_id, - status=payload.status, - trace_id=envelope.trace_id, - ) + await consume_with_failure_signal( + message, + domain="web_resume", + action="analyze", + envelope_type=Envelope[WebResumeAnalyzeRequest], + idempotency=self._idempotency, + publisher=self._publisher, + routing_key=self._callback_routing_key, + message_type="callback.analysis", + process=self._process, + failed_payload=self._failed_payload, + done_fields=analysis_done_fields, + expected_errors=(WebResumeAnalyzeError,), + ) - async def _run_and_build_payload( - self, - req: WebResumeAnalyzeRequest, - trace_id: str, + async def _process( + self, envelope: Envelope[WebResumeAnalyzeRequest] ) -> AnalysisCallbackPayload: - try: - result = await self._analyzer.analyze( - resume_id=req.resume_id, - url=req.url, - analyzed_document_id=req.analyzed_document_id, - ) - except WebResumeAnalyzeError as err: - log.warning( - "web_resume.analyze.domain_failed", - resume_id=req.resume_id, - code=err.code, - retriable=err.retriable, - trace_id=trace_id, - ) - return AnalysisCallbackPayload( - target_type="WEB", - target_id=req.resume_id, - status="FAILED", - error_code=err.code, - error_message=err.message, - retriable=err.retriable, - ) - except Exception as exc: - log.exception( - "web_resume.analyze.unexpected_failed", - resume_id=req.resume_id, - trace_id=trace_id, - ) - return AnalysisCallbackPayload( - target_type="WEB", - target_id=req.resume_id, - status="FAILED", - error_code="UNEXPECTED", - error_message=str(exc), - retriable=True, - ) - + req = envelope.payload + log.info( + "web_resume.analyze.start", + message_id=envelope.message_id, + resume_id=req.resume_id, + url=req.url, + trace_id=envelope.trace_id, + ) + result = await self._analyzer.analyze( + resume_id=req.resume_id, + url=req.url, + analyzed_document_id=req.analyzed_document_id, + ) return AnalysisCallbackPayload( target_type="WEB", target_id=req.resume_id, @@ -132,3 +78,13 @@ async def _run_and_build_payload( document_path=result.document_path, embedding_chunk_count=result.embedding_chunk_count, ) + + def _failed_payload( + self, req: WebResumeAnalyzeRequest, exc: Exception + ) -> AnalysisCallbackPayload: + return analysis_failed_payload( + target_type="WEB", + target_id=req.resume_id, + exc=exc, + domain_error=WebResumeAnalyzeError, + ) diff --git a/ai/src/ai_server/model/messages/feedback.py b/ai/src/ai_server/model/messages/feedback.py index bb2fbbe9..aef2c01c 100644 --- a/ai/src/ai_server/model/messages/feedback.py +++ b/ai/src/ai_server/model/messages/feedback.py @@ -70,6 +70,9 @@ class GenerateFeedbackRequest(BaseModel): # 직무 맞춤(JOB_TAILORED) 모드 전용. 회사명 + 채용공고(JD). '직무 적합도' 평가의 근거. target_company_name: str | None = None target_job_description: str | None = None + # 시도 상관관계(F5). Core 가 발행마다 발급 — 콜백에 그대로 에코해 Core 가 대체된 + # 이전 시도의 지연 FAILED 를 구분할 수 있게 한다. 구버전 요청은 None. + attempt_id: str | None = None class PanelBreakdownItem(BaseModel): @@ -126,3 +129,5 @@ class FeedbackCallbackPayload(BaseModel): error_code: str | None = None error_message: str | None = None retriable: bool | None = None + # 요청의 attemptId 에코(F5) — Core 가 stale FAILED(대체된 이전 시도)를 걸러내는 근거. + attempt_id: str | None = None diff --git a/ai/tests/test_cover_letter_consumer.py b/ai/tests/test_cover_letter_consumer.py index 58bd3e68..3671580d 100644 --- a/ai/tests/test_cover_letter_consumer.py +++ b/ai/tests/test_cover_letter_consumer.py @@ -105,3 +105,21 @@ async def test_domain_error_publishes_failed_cover_letter_callback() -> None: assert payload.status == "FAILED" assert payload.target_type == "COVER_LETTER" assert payload.error_code == "EMPTY_PDF_TEXT" + + +@pytest.mark.asyncio +async def test_unexpected_error_publishes_failed_callback() -> None: + """예상 못 한 예외도 FAILED 콜백(UNEXPECTED, retriable=true) — 공용 가드(F6) 회귀 고정.""" + analyzer = AsyncMock() + analyzer.analyze = AsyncMock(side_effect=RuntimeError("llm blew up")) + consumer, publisher = _make_consumer(analyzer) + + await consumer.handle(_incoming_message(_request_envelope())) + + publisher.publish.assert_awaited_once() + payload = publisher.publish.await_args.kwargs["payload"] + assert payload.status == "FAILED" + assert payload.error_code == "UNEXPECTED" + assert payload.error_message == "RuntimeError: llm blew up" + assert payload.retriable is True + assert payload.target_id == 9 diff --git a/ai/tests/test_feedback_consumer.py b/ai/tests/test_feedback_consumer.py index dff9c62c..0098e0bb 100644 --- a/ai/tests/test_feedback_consumer.py +++ b/ai/tests/test_feedback_consumer.py @@ -89,6 +89,7 @@ def _envelope( }, ], "contextDocumentIds": context_documents or [], + "attemptId": "att-1", }, "context": {"userId": 1, "sessionId": 50}, } @@ -136,6 +137,7 @@ async def test_consumer_generates_feedback_and_publishes_callback(): assert payload.session_id == 50 assert payload.overall_score == 85.0 assert payload.status == "OK" # 실패 신호 도입 후에도 성공 콜백은 OK(기본값) + assert payload.attempt_id == "att-1" # 요청의 attemptId 에코(F5) assert publisher.publish.await_args.kwargs["message_type"] == "callback.feedback" @@ -197,6 +199,7 @@ async def test_consumer_publishes_failed_callback_on_unexpected_error(): assert payload.error_code == "UNEXPECTED" assert payload.error_message == "RuntimeError: rag down" assert payload.retriable is True + assert payload.attempt_id == "att-1" # FAILED 도 에코 — Core 의 stale 판별 근거(F5) assert payload.overall_score is None assert publisher.publish.await_args.kwargs["message_type"] == "callback.feedback" diff --git a/ai/tests/test_resume_consumer.py b/ai/tests/test_resume_consumer.py index e1f5d306..ce4ee90d 100644 --- a/ai/tests/test_resume_consumer.py +++ b/ai/tests/test_resume_consumer.py @@ -174,3 +174,27 @@ async def test_parse_failure_raises_for_dlq() -> None: analyzer.analyze.assert_not_called() publisher.publish.assert_not_called() + + +@pytest.mark.asyncio +async def test_consumer_unmarks_when_success_publish_fails(): + """콜백 발행만 실패한 경우 — FAILED 오인 없이 원 예외로 DLQ + 멱등 unmark(F6). + unmark 가 없으면 재주입이 duplicate skip 으로 삼켜져 콜백이 영영 안 나간다.""" + analyzer = AsyncMock() + analyzer.analyze = AsyncMock( + return_value=MagicMock( + summary="s", + tech_stack=["Java"], + document_path="analyzed/resume/42/summary.md", + embedding_chunk_count=3, + ) + ) + store = LruIdempotencyStore(max_size=16) + consumer, publisher = _make_consumer(analyzer, idempotency=store) + publisher.publish = AsyncMock(side_effect=ConnectionError("mq down")) + + with pytest.raises(ConnectionError, match="mq down"): + await consumer.handle(_incoming_message(_request_envelope())) + + publisher.publish.assert_awaited_once() # FAILED 재발행 시도 없음 + assert store.is_seen_then_mark("req-1") is False diff --git a/ai/tests/test_tts_consumer.py b/ai/tests/test_tts_consumer.py index 1b13d44e..4a5ccda4 100644 --- a/ai/tests/test_tts_consumer.py +++ b/ai/tests/test_tts_consumer.py @@ -251,3 +251,57 @@ async def synthesize(self, text, *, voice): # 3문장 중 2번째 실패 → 성공 2개에 seq 0,1 연속 부여 assert [a["seq"] for a in notifier.audio] == [0, 1] assert publisher.published[0][2].status == "SUCCEEDED" + + +@pytest.mark.asyncio +async def test_tts_consumer_unmarks_when_publish_fails(): + """콜백 발행 실패 시 멱등 unmark 후 원 예외로 DLQ(F6) — 재주입 삼킴 방지.""" + from unittest.mock import AsyncMock, MagicMock + + from ai_server.messaging.idempotency import LruIdempotencyStore + + publisher = MagicMock() + publisher.publish = AsyncMock(side_effect=ConnectionError("mq down")) + store = LruIdempotencyStore(max_size=10) + consumer = TtsConsumer( + tts=MockTtsProvider(), + storage=_FakeStorage(), + publisher=publisher, + idempotency=store, + callback_routing_key="callback.tts", + voice="alloy", + key_template="interview/tts/{session_id}/{message_id}.mp3", + ) + + with pytest.raises(ConnectionError, match="mq down"): + await consumer.handle(_FakeMsg(_envelope_body())) + + assert store.is_seen_then_mark("m-1") is False + + +@pytest.mark.asyncio +async def test_tts_consumer_unmarks_when_prepublish_step_fails(): + """발행 이전 단계(키 조립 등)에서 죽어도 unmark 후 DLQ(F6) — + 마킹만 남긴 채 격리되면 재주입이 duplicate skip 으로 삼켜진다.""" + from unittest.mock import AsyncMock, MagicMock + + from ai_server.messaging.idempotency import LruIdempotencyStore + + publisher = MagicMock() + publisher.publish = AsyncMock() + store = LruIdempotencyStore(max_size=10) + consumer = TtsConsumer( + tts=MockTtsProvider(), + storage=_FakeStorage(), + publisher=publisher, + idempotency=store, + callback_routing_key="callback.tts", + voice="alloy", + key_template="interview/tts/{nope}.mp3", # 잘못된 템플릿 → KeyError + ) + + with pytest.raises(KeyError): + await consumer.handle(_FakeMsg(_envelope_body())) + + publisher.publish.assert_not_awaited() + assert store.is_seen_then_mark("m-1") is False diff --git a/ai/tests/test_voice_consumer.py b/ai/tests/test_voice_consumer.py index 0df1a2dc..bafdb969 100644 --- a/ai/tests/test_voice_consumer.py +++ b/ai/tests/test_voice_consumer.py @@ -227,3 +227,26 @@ async def test_consumer_idempotent_skip(): ) await consumer.handle(_StubMessage(_envelope())) publisher.publish.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_consumer_unmarks_when_publish_fails(): + """콜백 발행 실패 시 멱등 unmark 후 원 예외로 DLQ(F6) — + 재주입이 duplicate skip 으로 삼켜져 콜백이 영영 안 나가는 것을 방지.""" + publisher = MagicMock() + publisher.publish = AsyncMock(side_effect=ConnectionError("mq down")) + store = LruIdempotencyStore(max_size=10) + + consumer = VoiceConsumer( + stt=_stt_ok(), + storage=_storage_ok(), + publisher=publisher, + idempotency=store, + callback_routing_key="callback.voice", + filler_pattern=r"(?:음+|어+|그+|아+)", + ) + + with pytest.raises(ConnectionError, match="mq down"): + await consumer.handle(_StubMessage(_envelope())) + + assert store.is_seen_then_mark("voice-1") is False diff --git a/ai/tests/test_web_consumer.py b/ai/tests/test_web_consumer.py index 880f2f27..e338556a 100644 --- a/ai/tests/test_web_consumer.py +++ b/ai/tests/test_web_consumer.py @@ -106,3 +106,21 @@ async def test_domain_error_publishes_failed_callback() -> None: assert payload.error_code == "WEB_HTTP_STATUS" assert payload.retriable is False assert payload.target_type == "WEB" + + +@pytest.mark.asyncio +async def test_unexpected_error_publishes_failed_callback() -> None: + """도메인 에러 밖의 예상 못 한 예외도 FAILED 콜백으로 신호(UNEXPECTED, retriable=true) — + 공용 가드(F6) 전환으로 처리되는 경로의 회귀 고정.""" + analyzer = AsyncMock() + analyzer.analyze = AsyncMock(side_effect=RuntimeError("fetch blew up")) + consumer, publisher = _make_consumer(analyzer) + + await consumer.handle(_incoming_message(_request_envelope())) + + payload = _captured_payload(publisher) + assert payload.status == "FAILED" + assert payload.error_code == "UNEXPECTED" + assert payload.error_message == "RuntimeError: fetch blew up" + assert payload.retriable is True + assert payload.target_id == 11 From 5906f2c2287e1a69acb8c55a85f6f598dd09f7cf Mon Sep 17 00:00:00 2001 From: Jaeho Date: Sun, 23 Aug 2026 03:25:53 +0900 Subject: [PATCH 2/3] =?UTF-8?q?feat(backend):=20=ED=94=BC=EB=93=9C?= =?UTF-8?q?=EB=B0=B1=20=EC=8B=9C=EB=8F=84=20=EC=83=81=EA=B4=80=EA=B4=80?= =?UTF-8?q?=EA=B3=84(V30)=20+=20stale=20FAILED=20=EB=93=9C=EB=A1=AD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit generate.feedback 발행마다 attemptId(UUID)를 세션(feedback_attempt_id, V30)에 기록하고 payload 에 동봉 — 발행은 attemptId 커밋 이후(afterCommit 동기화)로 지연해 즉시 FAILED 가 이전 값과 대조돼 오드롭되는 경합을 막는다. 콜백의 attemptId 가 현재와 다르면(양쪽 non-null) FAILED 를 드롭해 대체된 이전 시도의 지연 실패가 마커를 되씌우지 않게 한다(null 은 구버전 호환 통과, 성공 콜백 미검사). 같은 계열의 POOL stale FAILED 도 풀 시딩 후엔 세션을 종료하지 않고 드롭. 테스트 +5. --- backend/CLAUDE.md | 5 +- .../application/FeedbackCallbackService.java | 16 +++++ .../application/QuestionsCallbackService.java | 7 +++ .../application/SessionFeedbackRequester.java | 39 +++++++++--- .../dto/FeedbackCallbackPayload.java | 6 +- .../dto/GenerateFeedbackPayload.java | 4 +- .../session/domain/InterviewSession.java | 9 +++ .../V30__add_session_feedback_attempt_id.sql | 6 ++ .../FeedbackCallbackServiceTest.java | 60 ++++++++++++++++++- .../QuestionsCallbackServiceTest.java | 21 +++++++ .../SessionFeedbackRequesterTest.java | 3 + 11 files changed, 162 insertions(+), 14 deletions(-) create mode 100644 backend/src/main/resources/db/migration/V30__add_session_feedback_attempt_id.sql diff --git a/backend/CLAUDE.md b/backend/CLAUDE.md index 35b01dec..c304ea0f 100644 --- a/backend/CLAUDE.md +++ b/backend/CLAUDE.md @@ -465,7 +465,10 @@ docker compose up -d SSE 는 휘발성이라 그 순간 미접속 클라이언트를 위해 `SessionFeedbackQueryService.get` 이 피드백 없음 + 마커 존재 시 404 `FEEDBACK_GENERATION_FAILED`(details.retriable)를 반환해 "생성 중"(FEEDBACK_NOT_READY)과 구분한다(`ownedFeedback` 공유 경로도 동일 분기). 마커는 - 성공 콜백·재생성 요청에서 클리어하며, 재생성의 `generate.feedback` 발행은 + 성공 콜백·재생성 요청에서 클리어하며, **시도 상관관계(V30)** 로 지연 FAILED 를 걸러낸다 — + `publishGenerateFeedback` 이 발행마다 새 attemptId(UUID)를 `feedback_attempt_id` 에 기록·동봉, + AI 가 콜백에 에코, `apply` 는 FAILED 콜백의 attemptId 가 현재와 다르면(양쪽 non-null) 드롭 + (null 은 구버전 호환으로 통과, 성공 콜백은 미검사). 재생성의 `generate.feedback` 발행은 `FeedbackRegenerateRequestedEvent` → AFTER_COMMIT 리스너로 — 마커 clear 커밋 전에 발행되는 역전(§"메시지 발행은 commit 이후" 규칙)을 막는다. AI 의 `errorMessage` 원문은 서버 로그에만 남기고 클라이언트에는 화이트리스트 문구만 보낸다 diff --git a/backend/src/main/java/com/stackup/stackup/session/application/FeedbackCallbackService.java b/backend/src/main/java/com/stackup/stackup/session/application/FeedbackCallbackService.java index 6638cd3c..78f9244a 100644 --- a/backend/src/main/java/com/stackup/stackup/session/application/FeedbackCallbackService.java +++ b/backend/src/main/java/com/stackup/stackup/session/application/FeedbackCallbackService.java @@ -75,6 +75,16 @@ public void apply(FeedbackCallbackEnvelope envelope) { return; } if (payload.isFailed()) { + // 대체된 이전 시도의 지연 FAILED(F5): attemptId 가 현재 시도와 다르면 마커를 + // 되씌우지 않고 드롭. null(구버전 AI/세션)은 검사 없이 통과 — 실패 알림 회귀 방지. + // 성공 콜백은 검사하지 않는다(늦은 성공 리포트도 유효, exists 가드가 중복 처리). + if (isStaleAttempt(session, payload)) { + log.info("callback.feedback stale FAILED from superseded attempt, drop. " + + "sessionId={}, payloadAttempt={}, currentAttempt={}", + sessionId, payload.attemptId(), session.getFeedbackAttemptId()); + markProcessed(envelope.messageId()); + return; + } applyFeedbackFailed(session, payload); markProcessed(envelope.messageId()); return; @@ -142,6 +152,12 @@ public record SessionFeedbackNotice(Long sessionId, Long feedbackId) { // 알린다 — 콜백 없이 DLQ 로만 가던 시절의 "피드백 생성 중 무기한 대기"를 끊는 게 목적. // errorMessage 는 str(exc) 원문(LLM 게이트웨이 주소·쿼터 상세 등)이라 서버 로그에만 남기고, // 클라이언트에는 화이트리스트 코드로 만든 안내 문구만 보낸다(QuestionsCallbackService 와 동일 원칙). + private static boolean isStaleAttempt(InterviewSession session, FeedbackCallbackPayload payload) { + return payload.attemptId() != null + && session.getFeedbackAttemptId() != null + && !payload.attemptId().equals(session.getFeedbackAttemptId()); + } + private void applyFeedbackFailed(InterviewSession session, FeedbackCallbackPayload payload) { log.warn("callback.feedback generation failed. sessionId={}, errorCode={}, retriable={}, message={}", session.getId(), payload.errorCode(), payload.retriable(), payload.errorMessage()); diff --git a/backend/src/main/java/com/stackup/stackup/session/application/QuestionsCallbackService.java b/backend/src/main/java/com/stackup/stackup/session/application/QuestionsCallbackService.java index bc18a9f3..b801c21e 100644 --- a/backend/src/main/java/com/stackup/stackup/session/application/QuestionsCallbackService.java +++ b/backend/src/main/java/com/stackup/stackup/session/application/QuestionsCallbackService.java @@ -170,6 +170,13 @@ private void applyInitialQuestion(InterviewSession session, QuestionsCallbackPay private void applyPoolFailed(InterviewSession session, QuestionsCallbackPayload payload) { log.warn("callback.questions POOL generation failed. sessionId={}, errorCode={}, retriable={}, message={}", session.getId(), payload.errorCode(), payload.retriable(), payload.errorMessage()); + // 재발행(이어하기 복구 등)으로 대체된 이전 시도의 지연 FAILED 가, 새 시도가 이미 풀을 + // 시딩한 건강한 세션을 종료시키지 않게 한다 (OK 분기의 중복 시딩 가드와 대칭). + if (poolRepository.countBySessionId(session.getId()) > 0) { + log.info("callback.questions POOL FAILED after pool already seeded — drop. sessionId={}", + session.getId()); + return; + } endSessionOnPoolFailure(session, "POOL_GENERATION_FAILED"); } diff --git a/backend/src/main/java/com/stackup/stackup/session/application/SessionFeedbackRequester.java b/backend/src/main/java/com/stackup/stackup/session/application/SessionFeedbackRequester.java index e359ba05..b9e4e4f1 100644 --- a/backend/src/main/java/com/stackup/stackup/session/application/SessionFeedbackRequester.java +++ b/backend/src/main/java/com/stackup/stackup/session/application/SessionFeedbackRequester.java @@ -30,6 +30,8 @@ import org.springframework.transaction.annotation.Transactional; import org.springframework.transaction.event.TransactionPhase; import org.springframework.transaction.event.TransactionalEventListener; +import org.springframework.transaction.support.TransactionSynchronization; +import org.springframework.transaction.support.TransactionSynchronizationManager; // 세션 COMPLETED commit 후 발화 → generate.feedback 발행 (US-24). // 멱등: session_feedbacks UNIQUE (session_id). 이미 피드백이 있으면 skip. @@ -92,6 +94,10 @@ public void publishGenerateFeedback(Long userId, Long sessionId, String reason) List voiceAnalyses = voiceAnalysisRepository.findByMessage_Session_Id(sessionId); + // 시도 상관관계(V30): 발행마다 새 attemptId 를 발급해 세션에 기록(dirty checking)하고 + // payload 에 동봉 — 재생성으로 대체된 이전 시도의 지연 FAILED 콜백을 걸러내는 근거. + String attemptId = java.util.UUID.randomUUID().toString(); + session.beginFeedbackAttempt(attemptId); GenerateFeedbackPayload payload = new GenerateFeedbackPayload( session.getId(), session.getMode().name(), @@ -104,16 +110,33 @@ public void publishGenerateFeedback(Long userId, Long sessionId, String reason) domainQuestionCounts, session.getTargetCompanyName(), session.getTargetJobDescription(), - selfIntroVoiceAnalysis(msgEntities, voiceAnalyses) + selfIntroVoiceAnalysis(msgEntities, voiceAnalyses), + attemptId ); - publisher.publishToAi( - properties.routingKeys().generateFeedback(), - payload, - new MessageContext(userId, session.getId(), null, null) - ); - log.info("generate.feedback published. sessionId={}, msgCount={}, ctx={}, reason={}", - session.getId(), messages.size(), contextDocumentIds.size(), reason); + // 발행은 attemptId 기록이 커밋된 뒤에만("메시지 발행은 commit 이후" 규칙) — 발행이 먼저면 + // AI 의 즉시 FAILED 콜백이 아직 이전 attemptId 가 커밋돼 있는 세션과 대조돼 stale 로 + // 오판·드롭되고, 발행 후 커밋 실패 시엔 진행 중 시도의 모든 콜백이 영원히 불일치한다. + Runnable send = () -> { + publisher.publishToAi( + properties.routingKeys().generateFeedback(), + payload, + new MessageContext(userId, session.getId(), null, null) + ); + log.info("generate.feedback published. sessionId={}, msgCount={}, ctx={}, reason={}, attemptId={}", + session.getId(), messages.size(), contextDocumentIds.size(), reason, attemptId); + }; + if (TransactionSynchronizationManager.isSynchronizationActive()) { + TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { + @Override + public void afterCommit() { + send.run(); + } + }); + } else { + // 트랜잭션 밖(단위 테스트 등) 폴백 — 즉시 발행. + send.run(); + } } private MessageItem toItem(InterviewMessage m) { diff --git a/backend/src/main/java/com/stackup/stackup/session/application/dto/FeedbackCallbackPayload.java b/backend/src/main/java/com/stackup/stackup/session/application/dto/FeedbackCallbackPayload.java index 0c866e73..69901e50 100644 --- a/backend/src/main/java/com/stackup/stackup/session/application/dto/FeedbackCallbackPayload.java +++ b/backend/src/main/java/com/stackup/stackup/session/application/dto/FeedbackCallbackPayload.java @@ -24,7 +24,9 @@ public record FeedbackCallbackPayload( String status, // OK(기본) | FAILED. null 이면 구버전 취급(OK). String errorCode, String errorMessage, - Boolean retriable + Boolean retriable, + // 요청 attemptId 에코(F5, V30). null(구버전 AI)이면 stale 검사를 건너뛴다. + String attemptId ) { // 구버전(실패 신호 없던 시절) 호출부·테스트 호환용 — status=OK 로 위임. public FeedbackCallbackPayload( @@ -36,7 +38,7 @@ public FeedbackCallbackPayload( ) { this(sessionId, overallScore, technicalAccuracy, logicScore, communicationScore, strengthsSummary, weaknessesSummary, improvementKeywords, studyPlan, highlights, - panelBreakdown, answerCoaching, reportS3Key, "OK", null, null, null); + panelBreakdown, answerCoaching, reportS3Key, "OK", null, null, null, null); } public boolean isFailed() { diff --git a/backend/src/main/java/com/stackup/stackup/session/application/dto/GenerateFeedbackPayload.java b/backend/src/main/java/com/stackup/stackup/session/application/dto/GenerateFeedbackPayload.java index 74d1d798..1ad65ed0 100644 --- a/backend/src/main/java/com/stackup/stackup/session/application/dto/GenerateFeedbackPayload.java +++ b/backend/src/main/java/com/stackup/stackup/session/application/dto/GenerateFeedbackPayload.java @@ -20,7 +20,9 @@ public record GenerateFeedbackPayload( String targetCompanyName, String targetJobDescription, // 자기소개 답변 단독의 음성 메트릭. 첫인상 평가에서 세션 평균 대신 사용. 없으면 null. - VoiceAnalysisSummary selfIntroVoiceAnalysis + VoiceAnalysisSummary selfIntroVoiceAnalysis, + // 시도 상관관계(F5). 발행마다 새 UUID — AI 가 콜백에 에코, Core 가 stale FAILED 를 걸러낸다. + String attemptId ) { public record MessageItem( diff --git a/backend/src/main/java/com/stackup/stackup/session/domain/InterviewSession.java b/backend/src/main/java/com/stackup/stackup/session/domain/InterviewSession.java index 7bd34fe0..28ca9147 100644 --- a/backend/src/main/java/com/stackup/stackup/session/domain/InterviewSession.java +++ b/backend/src/main/java/com/stackup/stackup/session/domain/InterviewSession.java @@ -124,6 +124,11 @@ public class InterviewSession extends BaseSoftDeleteEntity { @Column(name = "feedback_fail_retriable") private Boolean feedbackFailRetriable; + // 현재 진행 중인 피드백 생성 시도 ID(V30). 발행마다 새로 발급 — FAILED 콜백의 attemptId 가 + // 이 값과 다르면 대체된 이전 시도의 지연 신호로 보고 무시한다(실패 마커 재마킹 방지). + @Column(name = "feedback_attempt_id", length = 36) + private String feedbackAttemptId; + private InterviewSession(User user, String title, String memo, SessionMode mode, List jobCategories, Integer maxQuestions, Integer maxDurationMinutes, @@ -234,6 +239,10 @@ public void cancel() { this.status = SessionStatus.CANCELLED; } + public void beginFeedbackAttempt(String attemptId) { + this.feedbackAttemptId = attemptId; + } + public void markFeedbackFailed(Boolean retriable) { this.feedbackFailedAt = Instant.now(); this.feedbackFailRetriable = retriable; diff --git a/backend/src/main/resources/db/migration/V30__add_session_feedback_attempt_id.sql b/backend/src/main/resources/db/migration/V30__add_session_feedback_attempt_id.sql new file mode 100644 index 00000000..285bc685 --- /dev/null +++ b/backend/src/main/resources/db/migration/V30__add_session_feedback_attempt_id.sql @@ -0,0 +1,6 @@ +-- 피드백 생성 시도 상관관계(F5). V29 실패 마커는 "어느 시도의 실패인지"를 몰라서, +-- 재생성으로 대체된 이전 시도의 지연 FAILED 콜백이 새 시도가 진행 중인 세션에 실패 마커를 +-- 되씌울 수 있었다(멱등은 messageId 단위 — 시도 식별자가 payload 에 없음). +-- 발행마다 새 attemptId(UUID)를 발급해 여기에 기록하고 payload 로 왕복시켜, +-- FAILED 콜백의 attemptId 가 현재 값과 다르면 무시한다. +ALTER TABLE interview_sessions ADD COLUMN feedback_attempt_id VARCHAR(36); diff --git a/backend/src/test/java/com/stackup/stackup/session/application/FeedbackCallbackServiceTest.java b/backend/src/test/java/com/stackup/stackup/session/application/FeedbackCallbackServiceTest.java index 0788e429..6ad4c963 100644 --- a/backend/src/test/java/com/stackup/stackup/session/application/FeedbackCallbackServiceTest.java +++ b/backend/src/test/java/com/stackup/stackup/session/application/FeedbackCallbackServiceTest.java @@ -147,7 +147,7 @@ void apply_failedCallbackSkipsSaveAndPushesErrorSse() { FeedbackCallbackEnvelope env = envelope(50L, "fb-fail", new FeedbackCallbackPayload(50L, null, null, null, null, null, null, List.of(), List.of(), List.of(), List.of(), List.of(), null, - "FAILED", "UNEXPECTED", "boom: internal detail", true)); + "FAILED", "UNEXPECTED", "boom: internal detail", true, null)); when(processedMessageRepository.existsById("fb-fail")).thenReturn(false); when(sessionRepository.findById(50L)).thenReturn(Optional.of(session)); when(feedbackRepository.existsBySession_Id(50L)).thenReturn(false); @@ -189,6 +189,62 @@ void apply_successClearsFailureMarker() { verify(feedbackRepository).save(any(SessionFeedback.class)); } + @Test + void apply_staleFailedFromSupersededAttempt_isDropped() { + // 재생성으로 대체된 이전 시도의 지연 FAILED(F5) — 새 시도가 진행 중인 세션에 + // 실패 마커를 되씌우거나 실패 알림을 쏘면 안 된다. 멱등 마킹만 하고 드롭. + InterviewSession session = sessionFixture(50L); + session.beginFeedbackAttempt("att-2"); + FeedbackCallbackEnvelope env = envelope(50L, "fb-stale", + new FeedbackCallbackPayload(50L, null, null, null, null, null, null, + List.of(), List.of(), List.of(), List.of(), List.of(), null, + "FAILED", "UNEXPECTED", "late failure of attempt 1", true, "att-1")); + when(processedMessageRepository.existsById("fb-stale")).thenReturn(false); + when(sessionRepository.findById(50L)).thenReturn(Optional.of(session)); + when(feedbackRepository.existsBySession_Id(50L)).thenReturn(false); + + service.apply(env); + + assertThat(session.hasFeedbackFailure()).isFalse(); + verify(errorNotifier, never()).notify(any(), any(), any()); + verify(processedMessageRepository).save(any()); + } + + @Test + void apply_failedWithMatchingAttempt_marksFailure() { + InterviewSession session = sessionFixture(50L); + session.beginFeedbackAttempt("att-1"); + FeedbackCallbackEnvelope env = envelope(50L, "fb-match", + new FeedbackCallbackPayload(50L, null, null, null, null, null, null, + List.of(), List.of(), List.of(), List.of(), List.of(), null, + "FAILED", "UNEXPECTED", "boom", true, "att-1")); + when(processedMessageRepository.existsById("fb-match")).thenReturn(false); + when(sessionRepository.findById(50L)).thenReturn(Optional.of(session)); + when(feedbackRepository.existsBySession_Id(50L)).thenReturn(false); + + service.apply(env); + + assertThat(session.hasFeedbackFailure()).isTrue(); + } + + @Test + void apply_failedWithNullAttempt_stillMarks() { + // 구버전 AI(attemptId 미전송) 콜백은 검사 없이 통과 — 실패 알림 자체가 죽는 회귀 방지. + InterviewSession session = sessionFixture(50L); + session.beginFeedbackAttempt("att-1"); + FeedbackCallbackEnvelope env = envelope(50L, "fb-legacy-fail", + new FeedbackCallbackPayload(50L, null, null, null, null, null, null, + List.of(), List.of(), List.of(), List.of(), List.of(), null, + "FAILED", "UNEXPECTED", "boom", true, null)); + when(processedMessageRepository.existsById("fb-legacy-fail")).thenReturn(false); + when(sessionRepository.findById(50L)).thenReturn(Optional.of(session)); + when(feedbackRepository.existsBySession_Id(50L)).thenReturn(false); + + service.apply(env); + + assertThat(session.hasFeedbackFailure()).isTrue(); + } + @Test void apply_nullStatusTreatedAsOk() { // 구버전 콜백(status 미명시)은 OK 로 취급 — 기존 저장 경로 회귀 방지. @@ -196,7 +252,7 @@ void apply_nullStatusTreatedAsOk() { FeedbackCallbackEnvelope env = envelope(50L, "fb-legacy", new FeedbackCallbackPayload(50L, 80.0, null, null, null, null, null, List.of(), List.of(), List.of(), List.of(), List.of(), null, - null, null, null, null)); + null, null, null, null, null)); when(processedMessageRepository.existsById("fb-legacy")).thenReturn(false); when(sessionRepository.findById(50L)).thenReturn(Optional.of(session)); when(feedbackRepository.existsBySession_Id(50L)).thenReturn(false); diff --git a/backend/src/test/java/com/stackup/stackup/session/application/QuestionsCallbackServiceTest.java b/backend/src/test/java/com/stackup/stackup/session/application/QuestionsCallbackServiceTest.java index edd7ed2f..6d6d3b23 100644 --- a/backend/src/test/java/com/stackup/stackup/session/application/QuestionsCallbackServiceTest.java +++ b/backend/src/test/java/com/stackup/stackup/session/application/QuestionsCallbackServiceTest.java @@ -187,6 +187,27 @@ void apply_followupFailedWithoutPlaceholder_doesNotLeakInternalErrorMessage() { .doesNotContain("Errno"); } + @Test + void apply_poolFailedAfterPoolSeeded_isDropped() { + // 재발행으로 대체된 이전 시도의 지연 POOL FAILED — 새 시도가 이미 풀을 시딩했다면 + // 건강한 세션을 종료시키지 않고 드롭한다 (F5 리뷰 이관 저비용 가드). + InterviewSession session = sessionFixture(25L, SessionStatus.IN_PROGRESS); + QuestionsCallbackPayload payload = new QuestionsCallbackPayload( + 25L, "POOL", List.of(), null, null, null, null, null, null, + "FAILED", "GENERATION_FAILED", "late failure", true + ); + QuestionsCallbackEnvelope env = new QuestionsCallbackEnvelope( + "m-pool-stale", "callback.questions", "1", "t", null, "ai", payload, null); + when(processedMessageRepository.existsById("m-pool-stale")).thenReturn(false); + when(sessionRepository.findById(25L)).thenReturn(Optional.of(session)); + when(poolRepository.countBySessionId(25L)).thenReturn(3L); + + service.apply(env); + + assertThat(session.getStatus()).isEqualTo(SessionStatus.IN_PROGRESS); + verify(sessionRepository, never()).finishIfInProgress(any(), any(), any()); + } + // errorCode 가 아예 없어도(변형 producer·수동 DLQ 재주입) NPE 없이 폴백 코드로 처리한다 — // NPE → 롤백 → 재시도 루프는 실패 신호가 없애려던 무기한 대기의 재현이다. @Test diff --git a/backend/src/test/java/com/stackup/stackup/session/application/SessionFeedbackRequesterTest.java b/backend/src/test/java/com/stackup/stackup/session/application/SessionFeedbackRequesterTest.java index fb6301b3..3e48e606 100644 --- a/backend/src/test/java/com/stackup/stackup/session/application/SessionFeedbackRequesterTest.java +++ b/backend/src/test/java/com/stackup/stackup/session/application/SessionFeedbackRequesterTest.java @@ -75,6 +75,9 @@ void onSessionEnded_publishesGenerateFeedback() { GenerateFeedbackPayload payload = payloadCaptor.getValue(); assertThat(payload.messages()).hasSize(2); + // 시도 상관관계(F5): 발행마다 새 attemptId 를 발급해 세션에 기록하고 payload 에 동봉. + assertThat(payload.attemptId()).isNotBlank(); + assertThat(session.getFeedbackAttemptId()).isEqualTo(payload.attemptId()); VoiceAnalysisSummary summary = payload.voiceAnalysisSummary(); assertThat(summary.analyzedMessageCount()).isEqualTo(2); assertThat(summary.averageSpeakingRateWpm()).isEqualTo(150.0); From e4fd5cf6c412a1f30e059340b2252c7e9c301d46 Mon Sep 17 00:00:00 2001 From: Jaeho Date: Sun, 23 Aug 2026 03:26:02 +0900 Subject: [PATCH 3/3] =?UTF-8?q?docs:=20attemptId=20=EC=99=95=EB=B3=B5?= =?UTF-8?q?=C2=B7=EA=B0=80=EB=93=9C=20=ED=99=95=EC=9E=A5=20=EC=8A=A4?= =?UTF-8?q?=ED=8E=99=20=EB=B0=98=EC=98=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit messaging.md §5.10/§5.11 attemptId 필드·stale 드롭 규칙, §6 분석 4종 가드· voice/tts unmark_on_error(재주입 부작용 명시). database.md V30 컬럼. --- docs/database.md | 3 +++ docs/messaging.md | 20 +++++++++++++++++--- 2 files changed, 20 insertions(+), 3 deletions(-) diff --git a/docs/database.md b/docs/database.md index 62b6ca7e..9015048e 100644 --- a/docs/database.md +++ b/docs/database.md @@ -175,6 +175,9 @@ CREATE TABLE interview_sessions ( -- 새로고침한 클라이언트가 GET 피드백에서 "생성 중"(FEEDBACK_NOT_READY)과 "실패"(FEEDBACK_GENERATION_FAILED)를 구분하는 근거. feedback_failed_at TIMESTAMPTZ, feedback_fail_retriable BOOLEAN, + -- 피드백 생성 시도 ID(V30). 발행마다 새 UUID — 대체된 이전 시도의 지연 FAILED 콜백이 + -- 마커를 되씌우지 않게 콜백의 attemptId 에코와 대조한다 (messaging.md §5.10/§5.11). + feedback_attempt_id VARCHAR(36), status VARCHAR(20) NOT NULL DEFAULT 'READY' CHECK (status IN ('READY','IN_PROGRESS','INTERRUPTED','COMPLETED','CANCELLED')), total_question_count INT DEFAULT 0, diff --git a/docs/messaging.md b/docs/messaging.md index b2848ad8..0fa3e954 100644 --- a/docs/messaging.md +++ b/docs/messaging.md @@ -457,6 +457,7 @@ "sessionId": 99, "mode": "JOB_TAILORED", "jobCategory": "BACKEND", + "attemptId": "8b1f…-uuid", "messages": [ { "id": 1, "sequenceNumber": 1, "role": "INTERVIEWER", "content": "자기소개…", "category": "SELF_INTRODUCTION" }, { "id": 2, "sequenceNumber": 2, "role": "INTERVIEWEE", "content": "…", "parentMessageId": 1 } @@ -467,6 +468,9 @@ } ``` +- `attemptId`: 시도 상관관계(V30). Core 가 발행마다 새 UUID 를 발급해 + `interview_sessions.feedback_attempt_id` 에 기록하고 동봉 — AI 는 콜백에 그대로 에코한다. + ### 5.11 `callback.feedback` > `answerCoaching[]` 은 질문별 복기 — 답변(INTERVIEWEE) 메시지별 모범 답안·리라이트·한 줄 코칭(자기소개 제외). @@ -514,7 +518,8 @@ "status": "FAILED", "errorCode": "UNEXPECTED", "errorMessage": "...", - "retriable": true + "retriable": true, + "attemptId": "8b1f…-uuid" } } ``` @@ -529,6 +534,9 @@ 알리며, **실패 마커를 영속화**한다(`interview_sessions.feedback_failed_at`/`feedback_fail_retriable`, V29) — SSE 를 놓친 클라이언트도 GET 피드백의 404 `FEEDBACK_GENERATION_FAILED` 로 실패를 구분한다. 마커는 성공 콜백 도착·재생성 요청 시 클리어. `errorMessage` 원문은 서버 로그에만 남긴다(클라이언트 미노출). +- **시도 상관관계**: 콜백의 `attemptId`(요청 에코)가 세션의 현재 값과 다르면 **FAILED 콜백은 드롭** + 된다 — 재생성으로 대체된 이전 시도의 지연 실패가 새 시도의 마커를 되씌우지 않게. 어느 쪽이든 + null(구버전)이면 검사 없이 통과, 성공 콜백은 검사하지 않는다(중복은 `session_feedbacks` UNIQUE 가 처리). ### 5.12 `realtime.session.notify` ```json @@ -636,13 +644,19 @@ placeholder 를 `FAILED` 로 확정해 클라이언트의 턴이 잠기지 않 ### AI Server (aio-pika) - 컨슈머는 `async with message.process(requeue=False)` 패턴. - 도메인 예외 (`ResumeAnalyzeError` 등) 는 catch 하여 실패 callback 발행 (재시도 무의미). -- 생성 계열 3개(`questions`/`followup`/`feedback` consumer)는 공용 가드 +- 생성 3개(`questions`/`followup`/`feedback`) + 분석 4개(`resume`/`web`/`repository`/`cover_letter`) + consumer 는 공용 가드 (`messaging/consumers/failure_signal.py: consume_with_failure_signal`)를 쓴다 — envelope 파싱·멱등 체크 이후 **전 구간**(컨텍스트 빌드·진행 이벤트·생성·payload 조립)의 예외를 catch 해 항상 `status: FAILED` 콜백을 발행하고 ACK (세션·placeholder 가 "생성 중"에 무기한 멈추지 않게). 성공 콜백 발행 실패는 FAILED 오인 없이 원 예외로 DLQ(재처리 가능), 콜백을 하나도 못 낸 채 DLQ 로 가는 경로는 멱등 마킹을 해제(unmark)해 재주입이 삼켜지지 않게 한다. - errorMessage 는 `ExcType: msg` 형식 500자 상한. + errorMessage 는 `ExcType: msg` 형식 500자 상한. 분석 컨슈머의 도메인 에러 + (`ResumeAnalyzeError` 등)는 각 컨슈머의 FAILED payload 팩토리가 코드·retriable 분류를 보존한다. +- `voice`/`tts` consumer 는 실패 시나리오별 직접-발행 모델을 유지하되, 멱등 마킹 이후 전 구간을 + `unmark_on_error` 로 감싸 어떤 예외로든 콜백 없이 DLQ 로 가면 마킹을 해제한다(재주입 삼킴 방지). + **주의**: 재주입은 pre-publish 부수효과(TTS 합성·STT 재과금, 라이브 세그먼트 재전송, ai-log + 중복 행)를 재실행한다 — 운영자 수동 복구 수단으로만 쓴다. - 그 외 예외는 re-raise → nack(requeue=false) → DLX 로 routing. - 일시 장애의 in-process 재시도는 미구현 (Phase 2 — 아래 Quorum Queue 도입과 함께).