diff --git a/ai/CLAUDE.md b/ai/CLAUDE.md index 81392f87..452e9e77 100644 --- a/ai/CLAUDE.md +++ b/ai/CLAUDE.md @@ -267,7 +267,9 @@ uv run pytest -k "embedding" # 키워드 - 함수/변수: `snake_case` - 상수: `UPPER_SNAKE_CASE` - 클래스: `PascalCase` -- 타입 힌트 필수 (`from __future__ import annotations` 없이 PEP 604 union `int | None`) +- 타입 힌트 필수 — PEP 604 union(`int | None`) 사용. 파일 첫 줄의 + `from __future__ import annotations` 는 코드베이스 전반의 실제 컨벤션이므로 유지 + (과거 "없이" 라는 서술은 코드와 불일치했던 stale 규정 — 2026-08-23 정정) - async first — sync IO 사용 시 명시적 이유 상세 공통 규약: [`/docs/coding-conventions.md`](../docs/coding-conventions.md). @@ -385,31 +387,20 @@ docker run --env-file .env -p 8000:8000 stackup-ai 예상 못 한 예외로 죽어도(다른 4개 부가 평가는 각자 예외를 삼키는데 이것만 그러지 않으면 top-level `asyncio.gather` 가 통째로 취소된다) 빈 `FeedbackResult`로 대체해 피드백 발행 자체는 항상 이어지게 한다. -- **질문 풀/꼬리질문 생성 실패 신호 본 구현**: `questions_consumer`/`followup_consumer`가 메인 생성 - 호출(`generate()`/`stream()`)을 무방비로 두던 문제를 고쳤다 — 실패하면 예외가 그대로 새서 - DLQ로 격리되고 Core 는 아무 신호도 못 받아 세션이 "생성 중"에 무기한 멈췄다(꼬리질문은 Core 가 - 이미 선INSERT 한 placeholder 가 영원히 안 채워짐). 이제 두 consumer 모두 생성 호출을 try/except - 로 감싸 실패해도 항상 `QuestionPoolCallbackPayload`/`FollowupCallbackPayload`(`status=FAILED`, - `errorCode`, `errorMessage`, `retriable`)를 발행한다. `errorCode`는 `TypeError`(LLM 출력 스키마 - 불일치, 재시도 무의미 → `retriable=false`)와 그 외(`GENERATION_FAILED`, `retriable=true`)를 구분. - Core 쪽 처리는 [`backend/CLAUDE.md`](../backend/CLAUDE.md) 참고. 같은 김에 `questions_consumer`의 - 다문서 RAG 경로에도 `followup_rag_timeout_sec`와 대칭인 `questions_rag_timeout_sec`(기본 1.5s) - 하드 타임아웃을 추가하고, 모든 `ChatOpenAI` 호출에 `llm_pro_timeout_sec`(30s)/`llm_flash_timeout_sec` - (10s) 요청 타임아웃을 명시했다(이전엔 미설정 — SDK 기본값까지 무기한 대기 가능). - -- **피드백 생성 실패 신호 본 구현**: `feedback_consumer` 만 위 리팩터에서 빠져 있던 gap 을 닫았다 — - 패널·부가 평가는 내부 폴백(빈 결과/생략)으로 흡수되지만, 그 방어망 밖(트랜스크립트/RAG 컨텍스트 빌드, - payload 조립, 발행 등)의 예상 못 한 예외는 그대로 새서 DLQ 로만 격리되고 Core 는 아무 신호도 못 받아 - 세션이 "피드백 생성 중"에 무기한 멈췄다. `handle()` 본문을 payload 를 **반환**하는 `_process()` 로 - 추출하고 envelope 파싱·멱등 체크 이후의 생성 전 구간을 try/except 로 감싸, 실패 시 - `FeedbackCallbackPayload`(`status=FAILED`, `errorCode`, `errorMessage`(상한 500자), `retriable`)를 - 발행하고 ack 한다(`_publish_failed`). 분류는 questions/followup 과 동일 — `TypeError` 는 - `GENERATION_SCHEMA_INVALID`/`retriable=false`, 그 외 `UNEXPECTED`/`retriable=true`. **성공 콜백 - 발행 실패는 생성 실패가 아니다** — FAILED 오인 발행 없이 원 예외로 DLQ(변경 전과 동일, 재처리 - 가능). 콜백을 하나도 못 낸 채 DLQ 로 가는 경로(폴백 발행 실패 포함)는 `LruIdempotencyStore.unmark` - 로 마킹을 되돌려 재주입 시 duplicate skip 으로 삼켜지지 않게 한다. `status` 는 `GenerationStatus` - Literal 재사용, 기본값 `OK` 라 성공 콜백·구버전 소비자와 하위호환. - Core 쪽 처리는 [`backend/CLAUDE.md`](../backend/CLAUDE.md) 참고. +- **생성 실패 신호 공용 가드 본 구현 (F4 일원화)**: 질문 풀·꼬리질문·피드백 3개 consumer 의 실패 + 신호 메커니즘을 `messaging/consumers/failure_signal.py: consume_with_failure_signal` 하나로 + 통합했다. 이전엔 questions/followup 이 생성 호출만 try 로 감싸 컨텍스트 빌드·진행 이벤트·발행 + 구간이 무방비였고(같은 무기한 대기 버그를 세 번, 세 깊이로 수리), feedback 만 전 구간 가드였다. + 이제 3개 모두: envelope 파싱 실패는 raise→DLQ, 파싱·멱등 이후 **전 구간** 예외는 항상 + `status=FAILED` 콜백(+ack), 성공 콜백 발행 실패는 FAILED 오인 없이 원 예외로 DLQ(재처리 가능), + 콜백 0건 DLQ 경로는 `LruIdempotencyStore.unmark`. 각 consumer 는 `_process(envelope)`(성공 + 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 쪽 처리는 + [`backend/CLAUDE.md`](../backend/CLAUDE.md) 참고. (질문 풀 RAG `questions_rag_timeout_sec` + 1.5s 하드 타임아웃과 `llm_pro_timeout_sec` 30s/`llm_flash_timeout_sec` 10s 요청 타임아웃은 + 이전 리팩터에서 도입되어 유지된다.) - **질문 풀·피드백 생성 진행 이벤트 본 구현 (B2)**: 스트리밍이 없는 두 블로킹 생성 경로(질문 풀 Pro ≤30s, 피드백 병렬 gather ≈2분 예산)가 진행 중 무통보였던 것을 고쳤다. `SessionRealtimeNotifier.emit_progress` diff --git a/ai/src/ai_server/CLAUDE.md b/ai/src/ai_server/CLAUDE.md index da72c2de..0013f983 100644 --- a/ai/src/ai_server/CLAUDE.md +++ b/ai/src/ai_server/CLAUDE.md @@ -76,6 +76,10 @@ 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 절) ```python # messaging/consumers/resume_consumer.py (패턴) diff --git a/ai/src/ai_server/messaging/consumers/failure_signal.py b/ai/src/ai_server/messaging/consumers/failure_signal.py new file mode 100644 index 00000000..95fa585b --- /dev/null +++ b/ai/src/ai_server/messaging/consumers/failure_signal.py @@ -0,0 +1,125 @@ +from __future__ import annotations + +from collections.abc import Awaitable, Callable +from typing import Any, TypeVar + +import structlog +from aio_pika.abc import AbstractIncomingMessage +from pydantic import BaseModel + +from ai_server.messaging.idempotency import LruIdempotencyStore +from ai_server.messaging.publisher import CallbackPublisher +from ai_server.model.envelope import Envelope + +log = structlog.get_logger(__name__) + +ReqT = TypeVar("ReqT", bound=BaseModel) + + +def format_error_message(exc: Exception) -> str: + """str(exc) 는 LLM 응답 본문·입력 repr 까지 담길 수 있다 — 로그·와이어 크기 상한.""" + return f"{type(exc).__name__}: {exc}"[:500] + + +def classify_failure(exc: Exception, *, unexpected_code: str) -> tuple[str, bool]: + """(error_code, retriable) 분류의 단일 구현. TypeError(LLM 출력 스키마 불일치)는 + 같은 입력으로 재시도해도 똑같이 죽는다 → retriable=false. 그 외는 컨슈머별 계약 코드 + (questions/followup: GENERATION_FAILED, feedback: UNEXPECTED).""" + if isinstance(exc, TypeError): + return "GENERATION_SCHEMA_INVALID", False + return unexpected_code, True + + +async def consume_with_failure_signal( + message: AbstractIncomingMessage, + *, + domain: str, + envelope_type: type[Envelope[ReqT]], + idempotency: LruIdempotencyStore, + publisher: CallbackPublisher, + routing_key: str, + message_type: str, + process: Callable[[Envelope[ReqT]], Awaitable[BaseModel]], + failed_payload: Callable[[ReqT, Exception], BaseModel], + done_fields: Callable[[Any], dict[str, Any]] | None = None, +) -> None: + """생성 컨슈머(질문 풀·꼬리질문·피드백) 공용 실패 신호 가드. + + 같은 무기한 대기 버그를 컨슈머마다 다른 깊이로 세 번 고쳐 온 메커니즘의 단일 구현: + - envelope 파싱 실패는 재시도 무의미 → raise → DLQ (콜백 대상 식별 불가). + - `process` 의 전 구간(컨텍스트 빌드·진행 이벤트·생성·payload 조립) 예외는 항상 + FAILED 콜백으로 신호하고 ack — 콜백 없이 DLQ 로만 격리되면 Core 가 실패를 모른 채 + 세션이 "생성 중"에 무기한 멈춘다. + - 성공 payload 의 발행 실패는 생성 실패가 아니다 — FAILED 로 오인 발행하지 않고 + 원 예외로 DLQ 에 보내 재처리 가능하게 남긴다. + - 콜백을 하나도 못 낸 채 DLQ 로 가는 경로는 멱등 마킹을 되돌린다(unmark) — + 재주입 시 duplicate skip 으로 삼켜지지 않게. + """ + async with message.process(requeue=False): + try: + envelope = envelope_type.model_validate_json(message.body) + except Exception as exc: + log.error( + f"{domain}.parse.failed", + error=str(exc), + delivery_tag=message.delivery_tag, + ) + raise + + if idempotency.is_seen_then_mark(envelope.message_id): + log.info( + f"{domain}.idempotent.skip", + message_id=envelope.message_id, + trace_id=envelope.trace_id, + ) + return + + async def _publish(payload: BaseModel) -> None: + await publisher.publish( + routing_key=routing_key, + message_type=message_type, + payload=payload, + trace_id=envelope.trace_id, + correlation_id=envelope.message_id, + context=envelope.context, + ) + + session_id = getattr(envelope.payload, "session_id", None) + + 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, + ) + try: + # 팩토리 자체가 죽어도(검증 오류 등) 같은 안전망을 태운다 — + # try 밖이면 unmark 없이 DLQ 로 가서 재주입이 duplicate skip 으로 삼켜진다. + fallback = failed_payload(envelope.payload, exc) + await _publish(fallback) + 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, + ) + idempotency.unmark(envelope.message_id) + raise exc + return + + try: + await _publish(payload) + except Exception: + 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, + **(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 8f5f1a2a..0a01dc15 100644 --- a/ai/src/ai_server/messaging/consumers/feedback_consumer.py +++ b/ai/src/ai_server/messaging/consumers/feedback_consumer.py @@ -26,6 +26,11 @@ SelfIntroEvaluator, ) from ai_server.core.client import CoreClient +from ai_server.messaging.consumers.failure_signal import ( + classify_failure, + consume_with_failure_signal, + format_error_message, +) from ai_server.messaging.idempotency import LruIdempotencyStore from ai_server.messaging.publisher import CallbackPublisher from ai_server.messaging.session_notify import ( @@ -103,46 +108,17 @@ def __init__( self._session_notifier = session_notifier async def handle(self, message: AbstractIncomingMessage) -> None: - async with message.process(requeue=False): - try: - envelope = Envelope[GenerateFeedbackRequest].model_validate_json( - message.body - ) - except Exception as exc: - log.error( - "feedback.parse.failed", - error=str(exc), - delivery_tag=message.delivery_tag, - ) - raise - - if self._idempotency.is_seen_then_mark(envelope.message_id): - log.info( - "feedback.idempotent.skip", - message_id=envelope.message_id, - trace_id=envelope.trace_id, - ) - return - - try: - payload = await self._process(envelope) - except Exception as exc: # noqa: BLE001 - await self._publish_failed(envelope, exc) - return - - # 성공 payload 의 발행 실패는 생성 실패가 아니다 — FAILED 콜백으로 오인 발행하지 - # 않고 원 예외로 DLQ 에 보내 재처리 가능하게 남긴다(실패 신호 도입 전과 동일 동작). - try: - await self._publish_callback(envelope, payload) - except Exception: - self._idempotency.unmark(envelope.message_id) - raise - log.info( - "feedback.generate.done", - message_id=envelope.message_id, - session_id=envelope.payload.session_id, - trace_id=envelope.trace_id, - ) + await consume_with_failure_signal( + message, + domain="feedback", + envelope_type=Envelope[GenerateFeedbackRequest], + idempotency=self._idempotency, + publisher=self._publisher, + routing_key=self._callback_routing_key, + message_type="callback.feedback", + process=self._process, + failed_payload=self._failed_payload, + ) async def _process( self, envelope: Envelope[GenerateFeedbackRequest] @@ -256,56 +232,20 @@ async def _tracked(coro: Awaitable[T]) -> T: return payload - async def _publish_callback( - self, - envelope: Envelope[GenerateFeedbackRequest], - payload: FeedbackCallbackPayload, - ) -> None: - await self._publisher.publish( - routing_key=self._callback_routing_key, - message_type="callback.feedback", - payload=payload, - trace_id=envelope.trace_id, - correlation_id=envelope.message_id, - context=envelope.context, - ) - - async def _publish_failed( - self, envelope: Envelope[GenerateFeedbackRequest], exc: Exception - ) -> None: - """생성 중 예상 못 한 예외의 실패 신호. 콜백 없이 DLQ 로만 격리되면 Core 가 실패를 - 모른 채 세션이 '피드백 생성 중'에 무기한 멈춘다 — 항상 FAILED 콜백을 발행하고 - ack 한다. 폴백 발행마저 실패하면 멱등 마킹을 되돌리고 원 예외를 다시 던져 - DLQ 로 보낸다(최후 안전망 — 재주입 시 duplicate skip 으로 삼켜지지 않게).""" - req = envelope.payload - log.exception( - "feedback.generate.unexpected", - message_id=envelope.message_id, - session_id=req.session_id, - trace_id=envelope.trace_id, - ) - # questions/followup consumer 와 동일 분류 — TypeError(LLM 출력 스키마 불일치)는 - # 같은 입력으로 재시도해도 똑같이 죽는다 → retriable=false. - is_schema = isinstance(exc, TypeError) - payload = FeedbackCallbackPayload( + def _failed_payload( + self, req: GenerateFeedbackRequest, exc: Exception + ) -> FeedbackCallbackPayload: + """questions/followup consumer 와 동일 분류 — TypeError(LLM 출력 스키마 불일치)는 + 같은 입력으로 재시도해도 똑같이 죽는다 → retriable=false. 그 외는 UNEXPECTED + (패널·부가 평가 실패는 내부 폴백으로 흡수되므로 여기 걸리는 건 진짜 예상 밖 예외).""" + error_code, retriable = classify_failure(exc, unexpected_code="UNEXPECTED") + return FeedbackCallbackPayload( session_id=req.session_id, status="FAILED", - error_code="GENERATION_SCHEMA_INVALID" if is_schema else "UNEXPECTED", - # str(exc) 는 LLM 응답 본문·입력 repr 까지 담길 수 있다 — 로그·와이어 크기 상한. - error_message=f"{type(exc).__name__}: {exc}"[:500], - retriable=not is_schema, + error_code=error_code, + error_message=format_error_message(exc), + retriable=retriable, ) - try: - await self._publish_callback(envelope, payload) - except Exception: # noqa: BLE001 - log.exception( - "feedback.failed_callback.publish_failed", - message_id=envelope.message_id, - session_id=req.session_id, - trace_id=envelope.trace_id, - ) - self._idempotency.unmark(envelope.message_id) - raise exc async def _emit_progress( self, diff --git a/ai/src/ai_server/messaging/consumers/followup_consumer.py b/ai/src/ai_server/messaging/consumers/followup_consumer.py index 55dcf56f..e5a1ac2e 100644 --- a/ai/src/ai_server/messaging/consumers/followup_consumer.py +++ b/ai/src/ai_server/messaging/consumers/followup_consumer.py @@ -12,6 +12,11 @@ ) from ai_server.chain.sentence_split import next_sentences from ai_server.core.client import CoreClient +from ai_server.messaging.consumers.failure_signal import ( + classify_failure, + consume_with_failure_signal, + format_error_message, +) from ai_server.messaging.idempotency import LruIdempotencyStore from ai_server.messaging.publisher import CallbackPublisher from ai_server.messaging.session_notify import SessionRealtimeNotifier @@ -65,101 +70,44 @@ def __init__( self._rag_timeout_sec = rag_timeout_sec async def handle(self, message: AbstractIncomingMessage) -> None: - async with message.process(requeue=False): - try: - envelope = Envelope[GenerateFollowupRequest].model_validate_json( - message.body - ) - except Exception as exc: - log.error( - "followup.parse.failed", - error=str(exc), - delivery_tag=message.delivery_tag, - ) - raise - - if self._idempotency.is_seen_then_mark(envelope.message_id): - log.info( - "followup.idempotent.skip", - message_id=envelope.message_id, - trace_id=envelope.trace_id, - ) - return - - req = envelope.payload - log.info( - "followup.generate.start", - message_id=envelope.message_id, - session_id=req.session_id, - parent=req.parent_message_id, - trace_id=envelope.trace_id, - ) - - rag_context = await self._build_rag_context(req) - payload = await self._generate_followup_payload( - req, rag_context, trace_id=envelope.trace_id - ) - - await self._publisher.publish( - routing_key=self._callback_routing_key, - message_type="callback.questions", - payload=payload, - trace_id=envelope.trace_id, - correlation_id=envelope.message_id, - context=envelope.context, - ) - log.info( - "followup.generate.done", - message_id=envelope.message_id, - session_id=req.session_id, - status=payload.status, - trace_id=envelope.trace_id, - ) + await consume_with_failure_signal( + message, + domain="followup", + envelope_type=Envelope[GenerateFollowupRequest], + idempotency=self._idempotency, + publisher=self._publisher, + routing_key=self._callback_routing_key, + message_type="callback.questions", + process=self._process, + failed_payload=self._failed_payload, + done_fields=lambda p: {"status": p.status}, + ) - async def _generate_followup_payload( - self, - req: GenerateFollowupRequest, - rag_context: str, - *, - trace_id: str, + async def _process( + self, envelope: Envelope[GenerateFollowupRequest] ) -> FollowupCallbackPayload: - """꼬리질문 생성. Core 가 이미 선INSERT 한 placeholder 가 있으므로, 실패해도 - 콜백은 항상 나가야 placeholder 가 영원히 '생성 중'으로 남지 않는다.""" - try: - if self._streaming is not None and self._notifier is not None: - result = await self._stream_followup(req, rag_context, trace_id) - else: - result = await self._generator.generate( - job_category=req.job_category, - mode=req.mode, - previous_question=req.previous_question, - answer_text=req.answer_text, - context=rag_context, - parent_category=req.parent_category or "UNKNOWN", - expected_signal=req.parent_expected_signal or "(none)", - history=_format_history(req.history), - ) - except Exception as exc: # noqa: BLE001 - log.exception( - "followup.generate.failed", - session_id=req.session_id, - followup_message_id=req.followup_message_id, - trace_id=trace_id, - ) - return FollowupCallbackPayload( - session_id=req.session_id, - kind="FOLLOWUP", - parent_message_id=req.parent_message_id, - answer_message_id=req.answer_message_id, - followup_message_id=req.followup_message_id, - status="FAILED", - error_code=( - "GENERATION_SCHEMA_INVALID" - if isinstance(exc, TypeError) - else "GENERATION_FAILED" - ), - error_message=str(exc), - retriable=not isinstance(exc, TypeError), + req = envelope.payload + log.info( + "followup.generate.start", + message_id=envelope.message_id, + session_id=req.session_id, + parent=req.parent_message_id, + trace_id=envelope.trace_id, + ) + + rag_context = await self._build_rag_context(req) + if self._streaming is not None and self._notifier is not None: + result = await self._stream_followup(req, rag_context, envelope.trace_id) + else: + result = await self._generator.generate( + job_category=req.job_category, + mode=req.mode, + previous_question=req.previous_question, + answer_text=req.answer_text, + context=rag_context, + parent_category=req.parent_category or "UNKNOWN", + expected_signal=req.parent_expected_signal or "(none)", + history=_format_history(req.history), ) return FollowupCallbackPayload( session_id=req.session_id, @@ -172,6 +120,27 @@ async def _generate_followup_payload( followup_message_id=req.followup_message_id, ) + def _failed_payload( + self, req: GenerateFollowupRequest, exc: Exception + ) -> FollowupCallbackPayload: + """Core 가 이미 선INSERT 한 placeholder 가 있으므로, 실패해도 콜백은 항상 나가야 + placeholder 가 영원히 '생성 중'으로 남지 않는다. 상관 필드 3종은 전부 req + (envelope payload) 출처라 어느 단계에서 실패해도 조립 가능하다.""" + error_code, retriable = classify_failure( + exc, unexpected_code="GENERATION_FAILED" + ) + return FollowupCallbackPayload( + session_id=req.session_id, + kind="FOLLOWUP", + parent_message_id=req.parent_message_id, + answer_message_id=req.answer_message_id, + followup_message_id=req.followup_message_id, + status="FAILED", + error_code=error_code, + error_message=format_error_message(exc), + retriable=retriable, + ) + async def _stream_followup( self, req: GenerateFollowupRequest, rag_context: str, trace_id: str ) -> FollowupResult: diff --git a/ai/src/ai_server/messaging/consumers/questions_consumer.py b/ai/src/ai_server/messaging/consumers/questions_consumer.py index 414a8e71..294df353 100644 --- a/ai/src/ai_server/messaging/consumers/questions_consumer.py +++ b/ai/src/ai_server/messaging/consumers/questions_consumer.py @@ -7,6 +7,11 @@ from ai_server.chain.question_generation_chain import QuestionGenerator from ai_server.core.client import CoreClient +from ai_server.messaging.consumers.failure_signal import ( + classify_failure, + consume_with_failure_signal, + format_error_message, +) from ai_server.messaging.idempotency import LruIdempotencyStore from ai_server.messaging.publisher import CallbackPublisher from ai_server.messaging.session_notify import ( @@ -53,83 +58,94 @@ def __init__( self._session_notifier = session_notifier async def handle(self, message: AbstractIncomingMessage) -> None: - async with message.process(requeue=False): - try: - envelope = Envelope[GenerateQuestionsRequest].model_validate_json( - message.body - ) - except Exception as exc: - log.error( - "questions.parse.failed", - error=str(exc), - delivery_tag=message.delivery_tag, - ) - raise - - if self._idempotency.is_seen_then_mark(envelope.message_id): - log.info( - "questions.idempotent.skip", - message_id=envelope.message_id, - trace_id=envelope.trace_id, - ) - return + await consume_with_failure_signal( + message, + domain="questions", + envelope_type=Envelope[GenerateQuestionsRequest], + idempotency=self._idempotency, + publisher=self._publisher, + routing_key=self._callback_routing_key, + message_type="callback.questions", + process=self._process, + failed_payload=self._failed_payload, + done_fields=lambda p: { + "status": p.status, + "question_count": len(p.questions), + }, + ) - req = envelope.payload - effective_pool_size = max( - 1, - req.initial_question_count, - ) - log.info( - "questions.generate.start", - message_id=envelope.message_id, - session_id=req.session_id, - doc_count=len(req.documents), - max_questions=req.max_questions, - pool_size=effective_pool_size, - trace_id=envelope.trace_id, - ) + async def _process( + self, envelope: Envelope[GenerateQuestionsRequest] + ) -> QuestionPoolCallbackPayload: + req = envelope.payload + effective_pool_size = max( + 1, + req.initial_question_count, + ) + log.info( + "questions.generate.start", + message_id=envelope.message_id, + session_id=req.session_id, + doc_count=len(req.documents), + max_questions=req.max_questions, + pool_size=effective_pool_size, + trace_id=envelope.trace_id, + ) - await self._emit_progress( - session_id=req.session_id, - phase="CONTEXT_BUILDING", - message="면접 자료를 정리하고 있어요.", - trace_id=envelope.trace_id, - ) - context_text = await self._build_context(req) - await self._emit_progress( - session_id=req.session_id, - phase="GENERATING", - message="자료를 바탕으로 첫 질문을 만들고 있어요.", - trace_id=envelope.trace_id, - ) - payload = await self._generate_pool_payload( - req, context_text, effective_pool_size, trace_id=envelope.trace_id - ) - # 생성 실패(FAILED 콜백) 직전에 "마무리하고 있어요" 가 스치면 오해를 부른다 — 성공시에만. - if payload.status != "FAILED": - await self._emit_progress( - session_id=req.session_id, - phase="FINALIZING", - message="질문 준비를 마무리하고 있어요.", - trace_id=envelope.trace_id, - ) + await self._emit_progress( + session_id=req.session_id, + phase="CONTEXT_BUILDING", + message="면접 자료를 정리하고 있어요.", + trace_id=envelope.trace_id, + ) + context_text = await self._build_context(req) + await self._emit_progress( + session_id=req.session_id, + phase="GENERATING", + message="자료를 바탕으로 첫 질문을 만들고 있어요.", + trace_id=envelope.trace_id, + ) + pool = await self._generator.generate( + job_categories=req.job_categories, + mode=req.mode, + max_questions=effective_pool_size, + context=context_text, + recent_questions=req.recent_questions, + self_introduction=req.self_intro_answer, + target_company_name=req.target_company_name, + target_job_description=req.target_job_description, + focus_areas=req.focus_areas, + ) + # 실패는 예외로 위 가드에 넘어가므로 여기 도달 = 성공 — + # FAILED 콜백 직전에 "마무리하고 있어요" 가 스치는 오해가 구조적으로 없다. + await self._emit_progress( + session_id=req.session_id, + phase="FINALIZING", + message="질문 준비를 마무리하고 있어요.", + trace_id=envelope.trace_id, + ) + return QuestionPoolCallbackPayload( + session_id=req.session_id, + kind="POOL", + questions=pool.questions, + ) - await self._publisher.publish( - routing_key=self._callback_routing_key, - message_type="callback.questions", - payload=payload, - trace_id=envelope.trace_id, - correlation_id=envelope.message_id, - context=envelope.context, - ) - log.info( - "questions.generate.done", - message_id=envelope.message_id, - session_id=req.session_id, - status=payload.status, - question_count=len(payload.questions), - trace_id=envelope.trace_id, - ) + def _failed_payload( + self, req: GenerateQuestionsRequest, exc: Exception + ) -> QuestionPoolCallbackPayload: + """실패해도 콜백은 항상 나가야 세션 시작이 무기한 멈추지 않는다.""" + error_code, retriable = classify_failure( + exc, unexpected_code="GENERATION_FAILED" + ) + return QuestionPoolCallbackPayload( + session_id=req.session_id, + kind="POOL", + questions=[], + status="FAILED", + error_code=error_code, + error_message=format_error_message(exc), + retriable=retriable, + ) async def _emit_progress( self, *, session_id: int, phase: str, message: str, trace_id: str @@ -144,52 +160,6 @@ async def _emit_progress( trace_id=trace_id, ) - async def _generate_pool_payload( - self, - req: GenerateQuestionsRequest, - context_text: str, - effective_pool_size: int, - *, - trace_id: str, - ) -> QuestionPoolCallbackPayload: - """질문 풀 생성. 실패해도 콜백은 항상 나가야 세션 시작이 무기한 멈추지 않는다.""" - try: - pool = await self._generator.generate( - job_categories=req.job_categories, - mode=req.mode, - max_questions=effective_pool_size, - context=context_text, - recent_questions=req.recent_questions, - self_introduction=req.self_intro_answer, - target_company_name=req.target_company_name, - target_job_description=req.target_job_description, - focus_areas=req.focus_areas, - ) - except Exception as exc: # noqa: BLE001 - log.exception( - "questions.generate.failed", - session_id=req.session_id, - trace_id=trace_id, - ) - return QuestionPoolCallbackPayload( - session_id=req.session_id, - kind="POOL", - questions=[], - status="FAILED", - error_code=( - "GENERATION_SCHEMA_INVALID" - if isinstance(exc, TypeError) - else "GENERATION_FAILED" - ), - error_message=str(exc), - retriable=not isinstance(exc, TypeError), - ) - return QuestionPoolCallbackPayload( - session_id=req.session_id, - kind="POOL", - questions=pool.questions, - ) - async def _build_context(self, req: GenerateQuestionsRequest) -> str: base_context = _build_context(req.documents) if not self._core or not self._embedder: diff --git a/ai/tests/test_followup_consumer.py b/ai/tests/test_followup_consumer.py index f72c9b17..8f67756b 100644 --- a/ai/tests/test_followup_consumer.py +++ b/ai/tests/test_followup_consumer.py @@ -182,6 +182,85 @@ async def test_consumer_publishes_failed_followup_when_generate_raises(): assert payload.followup_message_id == 503 +@pytest.mark.asyncio +async def test_consumer_publishes_failed_followup_when_rag_build_raises(): + """생성 호출 밖(RAG 컨텍스트 빌드 등) 예외도 FAILED 콜백으로 신호 — + 예전엔 DLQ 로만 가서 placeholder 가 영원히 '생성 중'으로 남았다 (전 구간 가드, F4). + 상관 필드 3종은 envelope 출처라 어느 단계 실패든 채워져야 한다.""" + generator = MagicMock() + generator.generate = AsyncMock() + publisher = MagicMock() + publisher.publish = AsyncMock() + + consumer = FollowupConsumer( + generator=generator, + publisher=publisher, + idempotency=LruIdempotencyStore(max_size=10), + callback_routing_key="callback.questions", + ) + consumer._build_rag_context = AsyncMock(side_effect=RuntimeError("core down")) + + await consumer.handle(_StubMessage(_envelope())) # raise 없이 ack 경로 + + generator.generate.assert_not_awaited() + payload: FollowupCallbackPayload = publisher.publish.await_args.kwargs["payload"] + assert payload.status == "FAILED" + assert payload.error_code == "GENERATION_FAILED" + assert payload.error_message == "RuntimeError: core down" + assert payload.parent_message_id == 501 + assert payload.answer_message_id == 502 + assert payload.followup_message_id == 503 + + +@pytest.mark.asyncio +async def test_consumer_unmarks_and_reraises_when_failed_publish_also_fails(): + """폴백(FAILED 콜백) 발행마저 실패하면 멱등 unmark 후 원 예외로 DLQ.""" + generator = MagicMock() + generator.generate = AsyncMock(side_effect=RuntimeError("gateway 500")) + publisher = MagicMock() + publisher.publish = AsyncMock(side_effect=ConnectionError("mq down")) + store = LruIdempotencyStore(max_size=10) + + consumer = FollowupConsumer( + generator=generator, + publisher=publisher, + idempotency=store, + callback_routing_key="callback.questions", + ) + + with pytest.raises(RuntimeError, match="gateway 500"): + await consumer.handle(_StubMessage(_envelope())) + + assert store.is_seen_then_mark("m-1") is False + + +@pytest.mark.asyncio +async def test_consumer_does_not_send_failed_when_success_publish_fails(): + """생성 성공 후 콜백 발행만 실패 — FAILED 오인 없이 원 예외로 DLQ(unmark 포함).""" + generator = MagicMock() + generator.generate = AsyncMock( + return_value=FollowupResult(followup_question="추가 질문?") + ) + publisher = MagicMock() + publisher.publish = AsyncMock(side_effect=ConnectionError("mq down")) + store = LruIdempotencyStore(max_size=10) + + consumer = FollowupConsumer( + generator=generator, + publisher=publisher, + idempotency=store, + callback_routing_key="callback.questions", + ) + + with pytest.raises(ConnectionError, match="mq down"): + await consumer.handle(_StubMessage(_envelope())) + + publisher.publish.assert_awaited_once() # FAILED 재발행 시도 없음 + payload: FollowupCallbackPayload = publisher.publish.await_args.kwargs["payload"] + assert payload.status == "OK" + assert store.is_seen_then_mark("m-1") is False + + @pytest.mark.asyncio async def test_consumer_injects_followup_rag_context_when_available(): followup_result = FollowupResult( diff --git a/ai/tests/test_questions_consumer.py b/ai/tests/test_questions_consumer.py index 9c94a6e4..ac95474f 100644 --- a/ai/tests/test_questions_consumer.py +++ b/ai/tests/test_questions_consumer.py @@ -151,6 +151,125 @@ async def test_consumer_publishes_failed_pool_when_generate_raises(): assert payload.questions == [] +def _basic_body() -> bytes: + return _envelope( + { + "sessionId": 99, + "mode": "TECHNICAL", + "jobCategories": ["BACKEND"], + "documents": [], + "maxQuestions": 5, + "initialQuestionCount": 2, + } + ) + + +@pytest.mark.asyncio +async def test_consumer_publishes_failed_pool_when_context_build_raises(): + """생성 호출 밖(컨텍스트 빌드 등) 예외도 이제 FAILED 콜백으로 신호한다 — + 예전엔 reject→DLQ 로만 가서 세션 시작이 무기한 멈췄다 (전 구간 가드, F4).""" + generator = MagicMock() + generator.generate = AsyncMock() + publisher = MagicMock() + publisher.publish = AsyncMock() + + consumer = QuestionsConsumer( + generator=generator, + publisher=publisher, + idempotency=LruIdempotencyStore(max_size=10), + callback_routing_key="callback.questions", + ) + consumer._build_context = AsyncMock(side_effect=RuntimeError("core down")) + + await consumer.handle(_StubMessage(_basic_body())) # raise 없이 ack 경로 + + generator.generate.assert_not_awaited() + publisher.publish.assert_awaited_once() + payload: QuestionPoolCallbackPayload = publisher.publish.await_args.kwargs[ + "payload" + ] + assert payload.status == "FAILED" + assert payload.error_code == "GENERATION_FAILED" + assert payload.error_message == "RuntimeError: core down" + assert payload.retriable is True + + +@pytest.mark.asyncio +async def test_consumer_unmarks_and_reraises_when_failed_publish_also_fails(): + """폴백(FAILED 콜백) 발행마저 실패하면 멱등 unmark 후 원 예외로 DLQ — + 재주입 시 duplicate skip 으로 삼켜지지 않는다.""" + generator = MagicMock() + generator.generate = AsyncMock(side_effect=RuntimeError("gateway 500")) + publisher = MagicMock() + publisher.publish = AsyncMock(side_effect=ConnectionError("mq down")) + store = LruIdempotencyStore(max_size=10) + + consumer = QuestionsConsumer( + generator=generator, + publisher=publisher, + idempotency=store, + callback_routing_key="callback.questions", + ) + + with pytest.raises(RuntimeError, match="gateway 500"): + await consumer.handle(_StubMessage(_basic_body())) + + assert store.is_seen_then_mark("m-1") is False + + +@pytest.mark.asyncio +async def test_consumer_unmarks_when_failed_payload_factory_raises(): + """FAILED payload 팩토리 자체가 죽어도 같은 안전망 — unmark 후 원 예외로 DLQ. + (팩토리 호출이 가드 try 밖이면 마킹이 남아 재주입이 duplicate skip 으로 삼켜진다.)""" + generator = MagicMock() + generator.generate = AsyncMock(side_effect=RuntimeError("gateway 500")) + publisher = MagicMock() + publisher.publish = AsyncMock() + store = LruIdempotencyStore(max_size=10) + + consumer = QuestionsConsumer( + generator=generator, + publisher=publisher, + idempotency=store, + callback_routing_key="callback.questions", + ) + consumer._failed_payload = MagicMock(side_effect=ValueError("factory broken")) + + with pytest.raises(RuntimeError, match="gateway 500"): + await consumer.handle(_StubMessage(_basic_body())) + + publisher.publish.assert_not_awaited() + assert store.is_seen_then_mark("m-1") is False + + +@pytest.mark.asyncio +async def test_consumer_does_not_send_failed_when_success_publish_fails(): + """생성이 성공했는데 콜백 발행만 실패한 경우 — FAILED 로 오인 발행하지 않고 + 원 예외로 DLQ(멱등 unmark 포함, 재처리 가능).""" + generator = MagicMock() + generator.generate = AsyncMock(return_value=GeneratedQuestionPool(questions=[])) + publisher = MagicMock() + publisher.publish = AsyncMock(side_effect=ConnectionError("mq down")) + store = LruIdempotencyStore(max_size=10) + + consumer = QuestionsConsumer( + generator=generator, + publisher=publisher, + idempotency=store, + callback_routing_key="callback.questions", + ) + + with pytest.raises(ConnectionError, match="mq down"): + await consumer.handle(_StubMessage(_basic_body())) + + publisher.publish.assert_awaited_once() # FAILED 재발행 시도 없음 + payload: QuestionPoolCallbackPayload = publisher.publish.await_args.kwargs[ + "payload" + ] + assert payload.status == "OK" + assert store.is_seen_then_mark("m-1") is False + + @pytest.mark.asyncio async def test_consumer_marks_failure_non_retriable_on_schema_type_error(): """LLM 출력이 스키마와 안 맞아 TypeError 가 나면 재시도해도 같은 이유로 또 실패할 diff --git a/backend/CLAUDE.md b/backend/CLAUDE.md index 78a4ab6c..35b01dec 100644 --- a/backend/CLAUDE.md +++ b/backend/CLAUDE.md @@ -469,7 +469,9 @@ docker compose up -d `FeedbackRegenerateRequestedEvent` → AFTER_COMMIT 리스너로 — 마커 clear 커밋 전에 발행되는 역전(§"메시지 발행은 commit 이후" 규칙)을 막는다. AI 의 `errorMessage` 원문은 서버 로그에만 남기고 클라이언트에는 화이트리스트 문구만 보낸다 - (QuestionsCallbackService 와 동일 원칙). AI 쪽 발행은 [`ai/CLAUDE.md`](../ai/CLAUDE.md) 참고. + (QuestionsCallbackService 와 동일 원칙 — 이중 채널 발행 자체는 공용 + `common/sse/SessionErrorNotifier` + `SessionErrorNotice` 로 통합, F4). AI 쪽 발행은 + [`ai/CLAUDE.md`](../ai/CLAUDE.md) 참고. - **문장 단위 TTS 세그먼트 프록시 본 구현 (Part B)**: `InterviewMessageService.streamAudioSegment` + `GET /api/sessions/{sid}/messages/{mid}/audio/segments/{seq}?ext=`. AI 가 휘발성으로 쓴 라이브 세그먼트를 규칙(`interview/tts/{sid}/{mid}/seg-{seq}.{ext}`)으로 재구성해 프록시(DB 미기록). 소유권+ext 화이트리스트+seq>=0 검증으로 임의 키 노출 차단. - AI 호출 로깅 (US-30) 본 구현: `/api/internal/ai-logs` + `ai_request_logs` INSERT - **웹 이력서(URL) 본 구현 (US-09)**: `POST /api/resumes/web { url }`. AI 서버에 웹 분석이 이미 diff --git a/backend/src/main/java/com/stackup/stackup/common/sse/SessionErrorNotice.java b/backend/src/main/java/com/stackup/stackup/common/sse/SessionErrorNotice.java new file mode 100644 index 00000000..ab795b3d --- /dev/null +++ b/backend/src/main/java/com/stackup/stackup/common/sse/SessionErrorNotice.java @@ -0,0 +1,10 @@ +package com.stackup.stackup.common.sse; + +// SseEventType.ERROR 의 세션 도메인 payload (docs/event-stream.md §3.6). +// scope 로 실패 지점을 구분한다 (FOLLOWUP | FEEDBACK). +// errorMessage 가 아니라 message — AI 내부 원문이 아니라 사용자에게 보여줄 문구다 +// (원문 화이트리스트 책임은 각 콜백 서비스에 있다). +public record SessionErrorNotice( + Long sessionId, String scope, String errorCode, String message, Boolean retriable +) { +} diff --git a/backend/src/main/java/com/stackup/stackup/common/sse/SessionErrorNotifier.java b/backend/src/main/java/com/stackup/stackup/common/sse/SessionErrorNotifier.java new file mode 100644 index 00000000..1fbd53ef --- /dev/null +++ b/backend/src/main/java/com/stackup/stackup/common/sse/SessionErrorNotifier.java @@ -0,0 +1,21 @@ +package com.stackup.stackup.common.sse; + +import com.stackup.stackup.common.messaging.RealtimeNotifyEvent; +import lombok.RequiredArgsConstructor; +import org.springframework.context.ApplicationEventPublisher; +import org.springframework.stereotype.Component; + +// 세션 도메인 실패의 SSE ERROR 발행 단일 구현 — 세션 채널(라이브 화면)과 유저 채널(워크스페이스) +// 양쪽에 같은 notice 를 쏜다. Questions/Feedback 콜백 서비스가 각자 인라인으로 들고 있던 +// 이중 발행 복사본을 이곳으로 모았다. common → 도메인 역참조가 없도록 id 만 받는다. +@Component +@RequiredArgsConstructor +public class SessionErrorNotifier { + + private final ApplicationEventPublisher events; + + public void notify(Long sessionId, Long userId, SessionErrorNotice notice) { + events.publishEvent(RealtimeNotifyEvent.session(sessionId, SseEventType.ERROR, notice)); + events.publishEvent(RealtimeNotifyEvent.user(userId, SseEventType.ERROR, notice)); + } +} 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 c6a22b7b..6638cd3c 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 @@ -3,6 +3,8 @@ import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import com.stackup.stackup.common.exception.ApiErrorCode; +import com.stackup.stackup.common.sse.SessionErrorNotice; +import com.stackup.stackup.common.sse.SessionErrorNotifier; import com.stackup.stackup.common.messaging.domain.ProcessedMessage; import com.stackup.stackup.common.messaging.domain.ProcessedMessageRepository; import com.stackup.stackup.common.messaging.RealtimeNotifyEvent; @@ -41,6 +43,7 @@ public class FeedbackCallbackService { private final InterviewMessageRepository messageRepository; private final ProcessedMessageRepository processedMessageRepository; private final ApplicationEventPublisher events; + private final SessionErrorNotifier errorNotifier; @Transactional public void apply(FeedbackCallbackEnvelope envelope) { @@ -145,10 +148,8 @@ private void applyFeedbackFailed(InterviewSession session, FeedbackCallbackPaylo // 영속 마커(V29) — SSE ERROR 는 휘발성이라, 그 순간 미접속 클라이언트도 GET 피드백에서 // 실패를 구분할 수 있게 남긴다 (dirty checking 으로 UPDATE). session.markFeedbackFailed(payload.retriable()); - QuestionsCallbackService.SessionErrorNotice notice = new QuestionsCallbackService.SessionErrorNotice( - session.getId(), "FEEDBACK", FEEDBACK_FAILED_CODE, FEEDBACK_FAILED_MESSAGE, payload.retriable()); - events.publishEvent(RealtimeNotifyEvent.session(session.getId(), SseEventType.ERROR, notice)); - events.publishEvent(RealtimeNotifyEvent.user(session.getUser().getId(), SseEventType.ERROR, notice)); + errorNotifier.notify(session.getId(), session.getUser().getId(), new SessionErrorNotice( + session.getId(), "FEEDBACK", FEEDBACK_FAILED_CODE, FEEDBACK_FAILED_MESSAGE, payload.retriable())); } // REST(GET 피드백)와 같은 코드·문구 — 채널에 따라 다른 안내가 나가지 않게 단일 출처로 묶는다. 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 00ca717a..bc18a9f3 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 @@ -3,6 +3,8 @@ import com.stackup.stackup.common.messaging.domain.ProcessedMessage; import com.stackup.stackup.common.messaging.domain.ProcessedMessageRepository; import com.stackup.stackup.common.messaging.RealtimeNotifyEvent; +import com.stackup.stackup.common.sse.SessionErrorNotice; +import com.stackup.stackup.common.sse.SessionErrorNotifier; import com.stackup.stackup.common.sse.SseEventType; import com.stackup.stackup.session.application.dto.QuestionsCallbackEnvelope; import com.stackup.stackup.session.application.dto.QuestionsCallbackPayload; @@ -46,6 +48,7 @@ public class QuestionsCallbackService { private final SessionQuestionPoolRepository poolRepository; private final ProcessedMessageRepository processedMessageRepository; private final ApplicationEventPublisher events; + private final SessionErrorNotifier errorNotifier; @Transactional public void apply(QuestionsCallbackEnvelope envelope) { @@ -326,16 +329,15 @@ private void applyFollowupFailed(InterviewSession session, QuestionsCallbackPayl advanceToNextGeneral(session.getId()); } - // AI 가 보내는 errorMessage 는 `str(exc)` 그대로다 — LLM 게이트웨이 주소·조직 식별자·쿼터 - // 상세, 파싱 실패 시엔 모델 원문까지 들어올 수 있다. 이걸 SSE 로 흘리면 브라우저까지 그대로 - // 나간다(프론트는 쓰지도 않는다). 원문은 위 log.warn 이 서버에 남기고, 클라이언트에는 - // 우리가 정의한 코드로만 만든 안내 문구를 보낸다. + // AI 의 errorMessage 는 예외 원문 파생(`ExcType: msg`, AI 측 500자 상한 컨벤션)이라 LLM + // 게이트웨이 주소·조직 식별자·쿼터 상세가 들어올 수 있다 — 상한은 크기 제한일 뿐 새니타이즈가 + // 아니고 Core 가 강제하는 보장도 아니다. 이걸 SSE 로 흘리면 브라우저까지 그대로 나간다. + // 원문은 위 log.warn 이 서버에 남기고, 클라이언트에는 화이트리스트 코드로 만든 안내 문구만 + // 보낸다. 발행 자체는 공용 SessionErrorNotifier. private void publishErrorEvent(InterviewSession session, String scope, QuestionsCallbackPayload payload) { - SessionErrorNotice notice = new SessionErrorNotice( + errorNotifier.notify(session.getId(), session.getUser().getId(), new SessionErrorNotice( session.getId(), scope, safeErrorCode(payload.errorCode()), - userFacingMessage(payload.errorCode()), payload.retriable()); - events.publishEvent(RealtimeNotifyEvent.session(session.getId(), SseEventType.ERROR, notice)); - events.publishEvent(RealtimeNotifyEvent.user(session.getUser().getId(), SseEventType.ERROR, notice)); + userFacingMessage(payload.errorCode()), payload.retriable())); } public record SessionStateNotice(Long sessionId, String status, String reason) { @@ -385,17 +387,15 @@ public record SessionMessageNotice(Long sessionId, Long messageId, String reason ); private static final String UNKNOWN_ERROR_CODE = "GENERATION_FAILED"; + // code == null 가드 필수 — Map.of 는 containsKey(null)/getOrDefault(null,…) 에서 NPE 를 + // 던진다. errorCode 없는 FAILED 콜백(변형 producer·수동 DLQ 재주입)이 NPE → 롤백 → 재시도 + // 루프에 빠지면, 이 실패 신호가 없애려던 "생성 중" 무기한 대기가 그대로 재현된다. private static String safeErrorCode(String code) { - return ERROR_MESSAGES.containsKey(code) ? code : UNKNOWN_ERROR_CODE; + return code != null && ERROR_MESSAGES.containsKey(code) ? code : UNKNOWN_ERROR_CODE; } private static String userFacingMessage(String code) { - return ERROR_MESSAGES.getOrDefault(code, ERROR_MESSAGES.get(UNKNOWN_ERROR_CODE)); + return ERROR_MESSAGES.get(safeErrorCode(code)); } - // errorMessage 가 아니라 message — 내부 원문이 아니라 사용자에게 보여줄 문구다. - public record SessionErrorNotice( - Long sessionId, String scope, String errorCode, String message, Boolean retriable - ) { - } } diff --git a/backend/src/test/java/com/stackup/stackup/common/sse/SessionErrorNotifierTest.java b/backend/src/test/java/com/stackup/stackup/common/sse/SessionErrorNotifierTest.java new file mode 100644 index 00000000..15ab6300 --- /dev/null +++ b/backend/src/test/java/com/stackup/stackup/common/sse/SessionErrorNotifierTest.java @@ -0,0 +1,45 @@ +package com.stackup.stackup.common.sse; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +import com.stackup.stackup.common.messaging.RealtimeNotifyEvent; +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; + +@ExtendWith(MockitoExtension.class) +class SessionErrorNotifierTest { + + @Mock ApplicationEventPublisher events; + @InjectMocks SessionErrorNotifier notifier; + + @Test + void notify_publishesErrorToSessionAndUserChannels() { + SessionErrorNotice notice = new SessionErrorNotice( + 50L, "FEEDBACK", "FEEDBACK_GENERATION_FAILED", "피드백 생성에 실패했습니다.", true); + + notifier.notify(50L, 1L, notice); + + ArgumentCaptor cap = ArgumentCaptor.forClass(Object.class); + verify(events, times(2)).publishEvent(cap.capture()); + assertThat(cap.getAllValues()) + .allSatisfy(e -> { + RealtimeNotifyEvent rne = (RealtimeNotifyEvent) e; + assertThat(rne.type()).isEqualTo(SseEventType.ERROR); + assertThat(rne.payload()).isSameAs(notice); + }); + assertThat(cap.getAllValues()) + .extracting(e -> ((RealtimeNotifyEvent) e).channel()) + .containsExactlyInAnyOrder( + RealtimeNotifyEvent.Channel.SESSION, RealtimeNotifyEvent.Channel.USER); + assertThat(cap.getAllValues()) + .extracting(e -> ((RealtimeNotifyEvent) e).id()) + .containsExactlyInAnyOrder(50L, 1L); + } +} 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 3c1d9967..0788e429 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 @@ -2,6 +2,7 @@ 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.atLeastOnce; import static org.mockito.Mockito.never; import static org.mockito.Mockito.verify; @@ -9,6 +10,8 @@ import com.stackup.stackup.common.messaging.RealtimeNotifyEvent; import com.stackup.stackup.common.messaging.domain.ProcessedMessageRepository; +import com.stackup.stackup.common.sse.SessionErrorNotice; +import com.stackup.stackup.common.sse.SessionErrorNotifier; import com.stackup.stackup.common.sse.SseEventType; import com.stackup.stackup.session.application.dto.AnswerCoachingItem; import com.stackup.stackup.session.application.dto.FeedbackCallbackEnvelope; @@ -42,6 +45,7 @@ class FeedbackCallbackServiceTest { @Mock InterviewMessageRepository messageRepository; @Mock ProcessedMessageRepository processedMessageRepository; @Mock ApplicationEventPublisher events; + @Mock SessionErrorNotifier errorNotifier; @InjectMocks FeedbackCallbackService service; @Test @@ -151,25 +155,15 @@ void apply_failedCallbackSkipsSaveAndPushesErrorSse() { service.apply(env); verify(feedbackRepository, never()).save(any(SessionFeedback.class)); - ArgumentCaptor evCap = ArgumentCaptor.forClass(Object.class); - verify(events, atLeastOnce()).publishEvent(evCap.capture()); - List errors = evCap.getAllValues().stream() - .filter(RealtimeNotifyEvent.class::isInstance) - .map(RealtimeNotifyEvent.class::cast) - .filter(e -> e.type() == SseEventType.ERROR) - .toList(); - assertThat(errors).hasSize(2); // session + user 채널 - assertThat(errors).extracting(RealtimeNotifyEvent::channel) - .containsExactlyInAnyOrder(RealtimeNotifyEvent.Channel.SESSION, RealtimeNotifyEvent.Channel.USER); - assertThat(errors).allSatisfy(e -> { - QuestionsCallbackService.SessionErrorNotice notice = - (QuestionsCallbackService.SessionErrorNotice) e.payload(); - assertThat(notice.scope()).isEqualTo("FEEDBACK"); - assertThat(notice.errorCode()).isEqualTo("FEEDBACK_GENERATION_FAILED"); - // AI errorMessage 원문(내부 상세)은 클라이언트로 새지 않는다. - assertThat(notice.message()).doesNotContain("internal detail"); - assertThat(notice.retriable()).isTrue(); - }); + // 이중 채널 발행 자체는 SessionErrorNotifier(공용) 책임 — 여기서는 notice 내용만 검증. + ArgumentCaptor noticeCap = ArgumentCaptor.forClass(SessionErrorNotice.class); + verify(errorNotifier).notify(eq(50L), eq(1L), noticeCap.capture()); + SessionErrorNotice notice = noticeCap.getValue(); + assertThat(notice.scope()).isEqualTo("FEEDBACK"); + assertThat(notice.errorCode()).isEqualTo("FEEDBACK_GENERATION_FAILED"); + // AI errorMessage 원문(내부 상세)은 클라이언트로 새지 않는다. + assertThat(notice.message()).doesNotContain("internal detail"); + assertThat(notice.retriable()).isTrue(); verify(processedMessageRepository).save(any()); // 영속 마커(V29) — SSE 를 놓친 클라이언트가 GET 피드백으로 실패를 구분하는 근거. assertThat(session.hasFeedbackFailure()).isTrue(); 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 9e85aa31..edd7ed2f 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 @@ -2,6 +2,7 @@ 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.atLeastOnce; import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; @@ -9,6 +10,7 @@ import static org.mockito.Mockito.when; import com.stackup.stackup.common.messaging.RealtimeNotifyEvent; +import com.stackup.stackup.common.sse.SessionErrorNotice; import com.stackup.stackup.common.messaging.domain.ProcessedMessage; import com.stackup.stackup.common.messaging.domain.ProcessedMessageRepository; import com.stackup.stackup.common.sse.SseEventType; @@ -41,6 +43,7 @@ class QuestionsCallbackServiceTest { @Mock com.stackup.stackup.session.domain.SessionQuestionPoolRepository poolRepository; @Mock ProcessedMessageRepository processedMessageRepository; @Mock org.springframework.context.ApplicationEventPublisher events; + @Mock com.stackup.stackup.common.sse.SessionErrorNotifier errorNotifier; @InjectMocks QuestionsCallbackService service; @Test @@ -171,15 +174,40 @@ void apply_followupFailedWithoutPlaceholder_doesNotLeakInternalErrorMessage() { service.apply(env); - ArgumentCaptor ev = ArgumentCaptor.forClass(Object.class); - verify(events, atLeastOnce()).publishEvent(ev.capture()); - assertThat(ev.getAllValues()) - .filteredOn(e -> e instanceof RealtimeNotifyEvent) - .isNotEmpty() - .allSatisfy(e -> assertThat(e.toString()) - .doesNotContain("mindlogic") - .doesNotContain("org-abc123") - .doesNotContain("Errno")); + // 발행 메커니즘은 공용 SessionErrorNotifier 로 위임 — notice 내용으로 검증. + ArgumentCaptor cap = ArgumentCaptor.forClass(SessionErrorNotice.class); + verify(errorNotifier).notify(eq(21L), eq(1L), cap.capture()); + SessionErrorNotice notice = cap.getValue(); + assertThat(notice.scope()).isEqualTo("FOLLOWUP"); + assertThat(notice.errorCode()).isEqualTo("GENERATION_FAILED"); + assertThat(notice.retriable()).isTrue(); + assertThat(notice.toString()) + .doesNotContain("mindlogic") + .doesNotContain("org-abc123") + .doesNotContain("Errno"); + } + + // errorCode 가 아예 없어도(변형 producer·수동 DLQ 재주입) NPE 없이 폴백 코드로 처리한다 — + // NPE → 롤백 → 재시도 루프는 실패 신호가 없애려던 무기한 대기의 재현이다. + @Test + void apply_followupFailedWithNullErrorCode_fallsBackWithoutNpe() { + InterviewSession session = sessionFixture(24L, SessionStatus.IN_PROGRESS); + QuestionsCallbackPayload payload = new QuestionsCallbackPayload( + 24L, "FOLLOWUP", List.of(), null, null, null, null, null, null, + "FAILED", null, null, null + ); + QuestionsCallbackEnvelope env = new QuestionsCallbackEnvelope( + "m-fu-nullcode", "callback.questions", "1", "t", null, "ai", payload, null); + + when(processedMessageRepository.existsById("m-fu-nullcode")).thenReturn(false); + when(sessionRepository.findById(24L)).thenReturn(Optional.of(session)); + + service.apply(env); + + ArgumentCaptor cap = ArgumentCaptor.forClass(SessionErrorNotice.class); + verify(errorNotifier).notify(eq(24L), eq(1L), cap.capture()); + assertThat(cap.getValue().errorCode()).isEqualTo("GENERATION_FAILED"); + assertThat(cap.getValue().message()).isNotBlank(); } // 모르는 코드가 와도 그대로 실어 보내지 않는다 — errorCode 역시 AI 가 채우는 문자열이다. @@ -198,11 +226,11 @@ void apply_followupFailedWithUnknownErrorCode_fallsBackToKnownCode() { service.apply(env); - ArgumentCaptor ev = ArgumentCaptor.forClass(Object.class); - verify(events, atLeastOnce()).publishEvent(ev.capture()); - assertThat(ev.getAllValues()) - .filteredOn(e -> e instanceof RealtimeNotifyEvent) - .allSatisfy(e -> assertThat(e.toString()).doesNotContain("org-secret")); + ArgumentCaptor cap = ArgumentCaptor.forClass(SessionErrorNotice.class); + verify(errorNotifier).notify(eq(22L), eq(1L), cap.capture()); + // 모르는 코드는 화이트리스트 폴백 코드로 대체된다. + assertThat(cap.getValue().errorCode()).isEqualTo("GENERATION_FAILED"); + assertThat(cap.getValue().toString()).doesNotContain("org-secret"); } @Test @@ -454,7 +482,8 @@ void apply_followupFailed_marksPlaceholderFailedAndAdvances() { QuestionsCallbackPayload payload = new QuestionsCallbackPayload( 23L, "FOLLOWUP", null, 200L, null, null, null, "NORMAL", 304L, - "FAILED", "GENERATION_FAILED", "gateway timeout", true + "FAILED", "GENERATION_FAILED", + "RuntimeError: cannot connect to https://llm-gateway.internal (org-abc123)", true ); QuestionsCallbackEnvelope env = new QuestionsCallbackEnvelope( "m-ph-failed", "callback.questions", "1", "t", null, "ai", payload, null); @@ -478,6 +507,14 @@ void apply_followupFailed_marksPlaceholderFailedAndAdvances() { .isEqualTo(com.stackup.stackup.session.domain.MessageStatus.FAILED); // DONT_KNOW 와 동일하게 다음 일반질문으로 진행 — 풀이 비어 세션 종료(POOL_EXHAUSTED). assertThat(session.getStatus()).isEqualTo(SessionStatus.COMPLETED); + // 이 실패 경로에서 발행되는 어떤 이벤트에도 AI 내부 원문이 실리지 않는다 — + // notice 단건 검증(레거시 폴백 테스트)보다 넓은 전수 스캔을 유지한다. + ArgumentCaptor ev = ArgumentCaptor.forClass(Object.class); + verify(events, atLeastOnce()).publishEvent(ev.capture()); + assertThat(ev.getAllValues()) + .allSatisfy(e -> assertThat(e.toString()) + .doesNotContain("llm-gateway") + .doesNotContain("org-abc123")); } @Test diff --git a/docs/messaging.md b/docs/messaging.md index 4ac56729..b2848ad8 100644 --- a/docs/messaging.md +++ b/docs/messaging.md @@ -636,12 +636,13 @@ placeholder 를 `FAILED` 로 확정해 클라이언트의 턴이 잠기지 않 ### AI Server (aio-pika) - 컨슈머는 `async with message.process(requeue=False)` 패턴. - 도메인 예외 (`ResumeAnalyzeError` 등) 는 catch 하여 실패 callback 발행 (재시도 무의미). -- `questions_consumer`/`followup_consumer` 도 동일 패턴 — 생성 호출을 catch 해 항상 콜백을 발행한다 - (`QuestionPoolCallbackPayload`/`FollowupCallbackPayload` 의 `status: FAILED`). Core 가 이미 선INSERT 한 - 꼬리질문 placeholder 가 영원히 "생성 중"으로 남는 것을 방지. -- `feedback_consumer` 는 envelope 파싱·멱등 체크 이후 전 구간을 catch 해 예상 못 한 예외 시 - `FeedbackCallbackPayload` 의 `status: FAILED` 콜백을 발행하고 ACK 한다. 폴백 발행마저 실패하면 - 원 예외를 re-raise → DLQ (최후 안전망). +- 생성 계열 3개(`questions`/`followup`/`feedback` 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자 상한. - 그 외 예외는 re-raise → nack(requeue=false) → DLX 로 routing. - 일시 장애의 in-process 재시도는 미구현 (Phase 2 — 아래 Quorum Queue 도입과 함께).