From 35c45f8f6dc53bf43343956e0c88258292502f2c Mon Sep 17 00:00:00 2001 From: Jaeho Date: Sat, 22 Aug 2026 22:30:18 +0900 Subject: [PATCH 1/3] =?UTF-8?q?feat(ai):=20=EC=A7=88=EB=AC=B8=20=ED=92=80?= =?UTF-8?q?=C2=B7=ED=94=BC=EB=93=9C=EB=B0=B1=20=EC=83=9D=EC=84=B1=20?= =?UTF-8?q?=EC=A7=84=ED=96=89=20=EC=9D=B4=EB=B2=A4=ED=8A=B8=20=EB=B0=9C?= =?UTF-8?q?=ED=96=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - SessionRealtimeNotifier.emit_progress 추가 — realtime.session.notify 로 QUESTION_POOL_PROGRESS / FEEDBACK_PROGRESS 직접 발행 (휘발성, Core 미경유) - questions_consumer: CONTEXT_BUILDING→GENERATING→FINALIZING 순차 3단계 (생성 실패 시 FINALIZING 은 발행하지 않음 — FAILED 콜백 직전 오해 방지) - feedback_consumer: 세부 평가 5개가 gather 병렬이라 SCORING 완료 카운터 (completed/total) 방식, PREPARING/FINALIZING 전후 단계 포함 - runner 에서 followup 과 동일 notifier 인스턴스 재사용 주입 - 기존 AnalysisProgressNotifier(user 채널) 무변경 — 경로 불변 회귀 테스트로 고정 - 단위 테스트 8건 추가 (발행 시퀀스·직렬화·예외 삼킴·미주입 회귀) --- ai/src/ai_server/CLAUDE.md | 2 +- .../messaging/consumers/feedback_consumer.py | 104 +++++++++++++++--- .../messaging/consumers/questions_consumer.py | 39 +++++++ ai/src/ai_server/messaging/runner.py | 12 +- ai/src/ai_server/messaging/session_notify.py | 43 ++++++++ ai/src/ai_server/model/messages/realtime.py | 21 ++++ ai/tests/test_feedback_consumer.py | 58 ++++++++++ ai/tests/test_progress.py | 23 ++++ ai/tests/test_questions_consumer.py | 83 ++++++++++++++ ai/tests/test_session_notify.py | 64 +++++++++++ 10 files changed, 429 insertions(+), 20 deletions(-) diff --git a/ai/src/ai_server/CLAUDE.md b/ai/src/ai_server/CLAUDE.md index b23eeba0..da72c2de 100644 --- a/ai/src/ai_server/CLAUDE.md +++ b/ai/src/ai_server/CLAUDE.md @@ -74,7 +74,7 @@ class QuestionPoolCallbackPayload(BaseModel): (resume/repository/web/cover_letter/questions/followup/feedback/voice/tts) - 조립·기동은 `runner.py` 의 `MessagingRuntime` (§3), 연결은 `connection.py`, 콜백 발행은 `publisher.py`, 멱등은 `idempotency.py`(`LruIdempotencyStore`), - RealTime 직접 발행은 `progress.py`(분석 진행)·`session_notify.py`(델타/오디오) + RealTime 직접 발행은 `progress.py`(분석 진행, user 채널)·`session_notify.py`(델타/오디오/질문 풀·피드백 생성 진행, 세션 채널) - 모든 consumer는 envelope parsing → trace_context → 비즈니스 핸들러 호출 패턴 ```python diff --git a/ai/src/ai_server/messaging/consumers/feedback_consumer.py b/ai/src/ai_server/messaging/consumers/feedback_consumer.py index 46e0c467..a6b91cb6 100644 --- a/ai/src/ai_server/messaging/consumers/feedback_consumer.py +++ b/ai/src/ai_server/messaging/consumers/feedback_consumer.py @@ -2,6 +2,8 @@ import asyncio import re +from collections.abc import Awaitable +from typing import TypeVar import structlog from aio_pika.abc import AbstractIncomingMessage @@ -26,6 +28,10 @@ from ai_server.core.client import CoreClient from ai_server.messaging.idempotency import LruIdempotencyStore from ai_server.messaging.publisher import CallbackPublisher +from ai_server.messaging.session_notify import ( + FEEDBACK_PROGRESS_EVENT, + SessionRealtimeNotifier, +) from ai_server.model.envelope import Envelope from ai_server.model.messages.feedback import ( AnswerCoachingItem, @@ -39,6 +45,8 @@ log = structlog.get_logger(__name__) +T = TypeVar("T") + _SELF_INTRO_CATEGORY = "SELF_INTRODUCTION" _JOB_TAILORED_MODE = "JOB_TAILORED" _BEHAVIORAL_CATEGORY = "BEHAVIORAL" @@ -77,6 +85,7 @@ def __init__( answer_coach: AnswerCoach | None = None, coaching_max_answers: int = 30, coaching_concurrency: int = 5, + session_notifier: SessionRealtimeNotifier | None = None, ) -> None: self._generator = generator self._publisher = publisher @@ -91,6 +100,7 @@ def __init__( self._answer_coach = answer_coach self._coaching_max_answers = coaching_max_answers self._coaching_concurrency = max(1, coaching_concurrency) + self._session_notifier = session_notifier async def handle(self, message: AbstractIncomingMessage) -> None: async with message.process(requeue=False): @@ -124,6 +134,12 @@ async def handle(self, message: AbstractIncomingMessage) -> None: trace_id=envelope.trace_id, ) + await self._emit_progress( + session_id=req.session_id, + phase="PREPARING", + message="면접 기록을 정리하고 있어요.", + trace_id=envelope.trace_id, + ) transcript = _build_transcript(req.messages) score_basis = _build_score_basis(req.messages) rag_context = await self._build_rag_context(req) @@ -131,6 +147,34 @@ async def handle(self, message: AbstractIncomingMessage) -> None: req.voice_analysis_summary ) + # 세부 평가 5개가 병렬(gather)이라 순차 phase 로는 진행을 표현할 수 없다 — + # 각 태스크 완료 시점에 completed/total 카운터로 emit 한다. + scoring_total = 5 + scoring_done = 0 + + async def _tracked(coro: Awaitable[T]) -> T: + nonlocal scoring_done + task_result = await coro + scoring_done += 1 + await self._emit_progress( + session_id=req.session_id, + phase="SCORING", + message=f"세부 평가를 진행하고 있어요. ({scoring_done}/{scoring_total})", + trace_id=envelope.trace_id, + completed=scoring_done, + total=scoring_total, + ) + return task_result + + await self._emit_progress( + session_id=req.session_id, + phase="SCORING", + message="평가위원들이 답변을 검토하고 있어요.", + trace_id=envelope.trace_id, + completed=0, + total=scoring_total, + ) + # 종합 피드백 + 자기소개 첫인상 + 직무 적합도(직무 맞춤 모드)를 병렬 실행. # 첫인상·직무 적합도는 종합 점수(overall)에 미포함 — generator 가 모른 채 계산한 뒤 표시용으로 덧붙인다. ( @@ -140,22 +184,24 @@ async def handle(self, message: AbstractIncomingMessage) -> None: personality_item, answer_coaching, ) = await asyncio.gather( - self._generate_panel( - job_category=req.job_category, - mode=req.mode, - total_question_count=req.total_question_count, - end_reason=req.end_reason, - transcript=transcript, - score_basis=score_basis, - rag_context=rag_context, - voice_analysis_summary=voice_analysis_summary, - domain_question_counts=req.domain_question_counts, - session_id=req.session_id, + _tracked( + self._generate_panel( + job_category=req.job_category, + mode=req.mode, + total_question_count=req.total_question_count, + end_reason=req.end_reason, + transcript=transcript, + score_basis=score_basis, + rag_context=rag_context, + voice_analysis_summary=voice_analysis_summary, + domain_question_counts=req.domain_question_counts, + session_id=req.session_id, + ) ), - self._evaluate_self_intro(req, voice_analysis_summary), - self._evaluate_job_fit(req, transcript, rag_context), - self._evaluate_personality(req), - self._coach_answers(req), + _tracked(self._evaluate_self_intro(req, voice_analysis_summary)), + _tracked(self._evaluate_job_fit(req, transcript, rag_context)), + _tracked(self._evaluate_personality(req)), + _tracked(self._coach_answers(req)), ) # 빈 평가위원 항목(점수·내용 모두 없음)은 표시하지 않는다 — LLM 부분 응답이 빈 패널로 새는 것 방지. extras = [self_intro_item, *job_fit_items, personality_item] @@ -163,6 +209,12 @@ async def handle(self, message: AbstractIncomingMessage) -> None: e for e in extras if e is not None and _panel_has_content(e) ) + await self._emit_progress( + session_id=req.session_id, + phase="FINALIZING", + message="피드백 리포트를 정리하고 있어요.", + trace_id=envelope.trace_id, + ) payload = FeedbackCallbackPayload( session_id=req.session_id, overall_score=result.overall_score, @@ -194,6 +246,28 @@ async def handle(self, message: AbstractIncomingMessage) -> None: trace_id=envelope.trace_id, ) + async def _emit_progress( + self, + *, + session_id: int, + phase: str, + message: str, + trace_id: str, + completed: int | None = None, + total: int | None = None, + ) -> None: + if self._session_notifier is None: + return + await self._session_notifier.emit_progress( + event_type=FEEDBACK_PROGRESS_EVENT, + session_id=session_id, + phase=phase, + message=message, + trace_id=trace_id, + completed=completed, + total=total, + ) + async def _generate_panel( self, *, diff --git a/ai/src/ai_server/messaging/consumers/questions_consumer.py b/ai/src/ai_server/messaging/consumers/questions_consumer.py index e9f973e9..414a8e71 100644 --- a/ai/src/ai_server/messaging/consumers/questions_consumer.py +++ b/ai/src/ai_server/messaging/consumers/questions_consumer.py @@ -9,6 +9,10 @@ from ai_server.core.client import CoreClient from ai_server.messaging.idempotency import LruIdempotencyStore from ai_server.messaging.publisher import CallbackPublisher +from ai_server.messaging.session_notify import ( + QUESTION_POOL_PROGRESS_EVENT, + SessionRealtimeNotifier, +) from ai_server.model.envelope import Envelope from ai_server.model.messages.questions import ( DocumentContext, @@ -33,6 +37,7 @@ def __init__( embedder: EmbeddingProvider | None = None, rag_top_k: int = 5, rag_timeout_sec: float = 1.5, + session_notifier: SessionRealtimeNotifier | None = None, ) -> None: self._generator = generator self._publisher = publisher @@ -45,6 +50,7 @@ def __init__( self._embedder = embedder self._rag_top_k = rag_top_k self._rag_timeout_sec = rag_timeout_sec + self._session_notifier = session_notifier async def handle(self, message: AbstractIncomingMessage) -> None: async with message.process(requeue=False): @@ -83,10 +89,30 @@ async def handle(self, message: AbstractIncomingMessage) -> None: 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._publisher.publish( routing_key=self._callback_routing_key, @@ -105,6 +131,19 @@ async def handle(self, message: AbstractIncomingMessage) -> None: trace_id=envelope.trace_id, ) + async def _emit_progress( + self, *, session_id: int, phase: str, message: str, trace_id: str + ) -> None: + if self._session_notifier is None: + return + await self._session_notifier.emit_progress( + event_type=QUESTION_POOL_PROGRESS_EVENT, + session_id=session_id, + phase=phase, + message=message, + trace_id=trace_id, + ) + async def _generate_pool_payload( self, req: GenerateQuestionsRequest, diff --git a/ai/src/ai_server/messaging/runner.py b/ai/src/ai_server/messaging/runner.py index 0060c338..4f7b2587 100644 --- a/ai/src/ai_server/messaging/runner.py +++ b/ai/src/ai_server/messaging/runner.py @@ -197,6 +197,12 @@ def __init__(self, settings: Settings) -> None: progress_notifier=self._progress_notifier, ) + # 세션 채널 휘발성 이벤트 발행기 — 꼬리질문 델타/TTS 와 질문 풀·피드백 진행 이벤트가 공유. + session_notifier = SessionRealtimeNotifier( + publisher=self._realtime_publisher, + routing_key="realtime.session.notify", + ) + # 질문 풀 생성 (US-18) question_generator = LlmQuestionGenerator( build_question_generation_chain(settings, core_client=core_client) @@ -210,6 +216,7 @@ def __init__(self, settings: Settings) -> None: core_client=core_client, embedder=embedder, rag_timeout_sec=settings.questions_rag_timeout_sec, + session_notifier=session_notifier, ) # 꼬리질문 생성 (US-19) @@ -219,10 +226,6 @@ def __init__(self, settings: Settings) -> None: streaming_followup_generator = build_streaming_followup_generator( settings, core_client=core_client ) - session_notifier = SessionRealtimeNotifier( - publisher=self._realtime_publisher, - routing_key="realtime.session.notify", - ) # TTS provider 는 꼬리질문 인라인 세그먼트 합성과 질문 TTS 양쪽에서 재사용한다. tts = build_tts_provider(settings) self._followup_consumer = FollowupConsumer( @@ -272,6 +275,7 @@ def __init__(self, settings: Settings) -> None: answer_coach=LlmAnswerCoach( build_answer_coaching_chain(settings, core_client=core_client) ), + session_notifier=session_notifier, ) # 음성 답변 STT + 분석 (Phase 2) diff --git a/ai/src/ai_server/messaging/session_notify.py b/ai/src/ai_server/messaging/session_notify.py index cd82cb57..0c31656d 100644 --- a/ai/src/ai_server/messaging/session_notify.py +++ b/ai/src/ai_server/messaging/session_notify.py @@ -9,12 +9,16 @@ SessionMessageAudioData, SessionMessageDeltaData, SessionNotifyPayload, + SessionProgressData, + SessionProgressNotifyPayload, ) log = structlog.get_logger(__name__) SESSION_MESSAGE_DELTA_EVENT = "SESSION_MESSAGE_DELTA" SESSION_MESSAGE_AUDIO_EVENT = "SESSION_MESSAGE_AUDIO" +QUESTION_POOL_PROGRESS_EVENT = "QUESTION_POOL_PROGRESS" +FEEDBACK_PROGRESS_EVENT = "FEEDBACK_PROGRESS" class SessionRealtimeNotifier: @@ -91,3 +95,42 @@ async def emit_audio( seq=seq, trace_id=trace_id, ) + + async def emit_progress( + self, + *, + event_type: str, + session_id: int, + phase: str, + message: str, + trace_id: str, + completed: int | None = None, + total: int | None = None, + ) -> None: + payload = SessionProgressNotifyPayload( + event_type=event_type, + data=SessionProgressData( + session_id=session_id, + phase=phase, + message=message, + completed=completed, + total=total, + ), + ) + try: + await self._publisher.publish( + routing_key=self._routing_key, + message_type=self._routing_key, + payload=payload, + trace_id=trace_id, + correlation_id=f"progress-{session_id}-{phase}-{completed or 0}", + context=MessageContext(session_id=session_id), + ) + except Exception: + log.warning( + "session.progress.publish_failed", + session_id=session_id, + event_type=event_type, + phase=phase, + trace_id=trace_id, + ) diff --git a/ai/src/ai_server/model/messages/realtime.py b/ai/src/ai_server/model/messages/realtime.py index 61e64ddb..9d8add24 100644 --- a/ai/src/ai_server/model/messages/realtime.py +++ b/ai/src/ai_server/model/messages/realtime.py @@ -39,6 +39,27 @@ class SessionNotifyPayload(BaseModel): data: SessionMessageDeltaData +# AI -> RealTime 세션 채널 직접 발행. 질문 풀/피드백 생성 진행 상황(휘발성). +# eventType: QUESTION_POOL_PROGRESS | FEEDBACK_PROGRESS +class SessionProgressData(BaseModel): + model_config = camel_config() + + session_id: int + # 질문 풀: CONTEXT_BUILDING | GENERATING | FINALIZING + # 피드백: PREPARING | SCORING | FINALIZING + phase: str + message: str + completed: int | None = None # 피드백 SCORING 에서만: 완료된 세부 평가 수 + total: int | None = None + + +class SessionProgressNotifyPayload(BaseModel): + model_config = camel_config() + + event_type: str + data: SessionProgressData + + # AI -> RealTime 세션 채널 직접 발행. 꼬리질문 문장 단위 TTS 세그먼트(휘발성). class SessionMessageAudioData(BaseModel): model_config = camel_config() diff --git a/ai/tests/test_feedback_consumer.py b/ai/tests/test_feedback_consumer.py index 9da41207..c66d8aae 100644 --- a/ai/tests/test_feedback_consumer.py +++ b/ai/tests/test_feedback_consumer.py @@ -996,3 +996,61 @@ async def test_consumer_skips_personality_in_technical_mode(): evaluator.evaluate.assert_not_awaited() payload: FeedbackCallbackPayload = publisher.publish.await_args.kwargs["payload"] assert all(b.evaluator != "인성·자소서" for b in payload.panel_breakdown) + + +@pytest.mark.asyncio +async def test_consumer_emits_feedback_progress_with_scoring_counter(): + """B2: 피드백 생성 중 세션 채널 진행 이벤트 — PREPARING → SCORING(0/5) → + 세부 평가 완료마다 카운터 증가(1..5) → FINALIZING 순서·개수를 고정한다.""" + generator = _generator() + publisher = MagicMock() + publisher.publish = AsyncMock() + session_notifier = MagicMock() + session_notifier.emit_progress = AsyncMock() + + consumer = FeedbackConsumer( + generator=generator, + publisher=publisher, + idempotency=LruIdempotencyStore(max_size=10), + callback_routing_key="callback.feedback", + core_client=MagicMock(), + embedder=None, + session_notifier=session_notifier, + ) + await consumer.handle(_StubMessage(_envelope())) + + calls = session_notifier.emit_progress.await_args_list + assert len(calls) == 8 # 1 PREPARING + 1 SCORING(0/5) + 5 완료 + 1 FINALIZING + phases = [c.kwargs["phase"] for c in calls] + assert phases[0] == "PREPARING" + assert phases[1] == "SCORING" + assert (calls[1].kwargs["completed"], calls[1].kwargs["total"]) == (0, 5) + assert phases[-1] == "FINALIZING" + completions = [ + c.kwargs["completed"] + for c in calls + if c.kwargs["phase"] == "SCORING" and c.kwargs["completed"] + ] + assert sorted(completions) == [1, 2, 3, 4, 5] + assert all(c.kwargs["event_type"] == "FEEDBACK_PROGRESS" for c in calls) + assert all(c.kwargs["session_id"] == 50 for c in calls) + publisher.publish.assert_awaited_once() + + +@pytest.mark.asyncio +async def test_consumer_without_session_notifier_still_publishes_callback(): + """B2 회귀 고정: notifier 미주입(None) 이면 진행 이벤트 없이 기존 흐름 그대로.""" + generator = _generator() + publisher = MagicMock() + publisher.publish = AsyncMock() + + consumer = FeedbackConsumer( + generator=generator, + publisher=publisher, + idempotency=LruIdempotencyStore(max_size=10), + callback_routing_key="callback.feedback", + core_client=MagicMock(), + embedder=None, + ) + await consumer.handle(_StubMessage(_envelope())) + publisher.publish.assert_awaited_once() diff --git a/ai/tests/test_progress.py b/ai/tests/test_progress.py index 17ea0882..5c653313 100644 --- a/ai/tests/test_progress.py +++ b/ai/tests/test_progress.py @@ -143,3 +143,26 @@ async def progress(phase: str, _message: str) -> None: ) assert phases == ["EXTRACTING", "SUMMARIZING", "EMBEDDING"] + + +@pytest.mark.asyncio +async def test_user_channel_emit_unchanged_by_session_progress_addition() -> None: + """B2 회귀 고정: 세션 채널 진행 이벤트(emit_progress) 도입 후에도 분석 진행은 + user 채널(realtime.user.notify) + user_id 컨텍스트 경로 그대로여야 한다.""" + notifier, publisher = _notifier() + + await notifier.emit( + user_id=5, + target_type="RESUME", + target_id=1, + phase="EXTRACTING", + message="추출 중", + trace_id="t-1", + ) + + kwargs = publisher.publish.await_args.kwargs + assert kwargs["routing_key"] == "realtime.user.notify" + assert kwargs["message_type"] == "realtime.user.notify" + assert kwargs["context"].user_id == 5 + assert kwargs["context"].session_id is None # 세션 채널로 새지 않는다 + assert kwargs["payload"].event_type == ANALYSIS_PROGRESS_EVENT diff --git a/ai/tests/test_questions_consumer.py b/ai/tests/test_questions_consumer.py index d085a9b9..9c94a6e4 100644 --- a/ai/tests/test_questions_consumer.py +++ b/ai/tests/test_questions_consumer.py @@ -588,3 +588,86 @@ def test_generated_question_parses_evidence_and_signal(): q2 = GeneratedQuestion.model_validate({"category": "BEHAVIORAL", "question": "Q?"}) assert q2.target_evidence == "" assert q2.expected_signal == "" + + +@pytest.mark.asyncio +async def test_consumer_emits_pool_progress_phases_in_order(): + """B2: 질문 풀 생성 중 세션 채널 진행 이벤트가 3단계 순서대로 발행돼야 한다.""" + generator = MagicMock() + generator.generate = AsyncMock( + return_value=GeneratedQuestionPool( + questions=[GeneratedQuestion(category="CS_FUNDAMENTAL", question="q?")] + ) + ) + publisher = MagicMock() + publisher.publish = AsyncMock() + session_notifier = MagicMock() + session_notifier.emit_progress = AsyncMock() + + consumer = QuestionsConsumer( + generator=generator, + publisher=publisher, + idempotency=LruIdempotencyStore(max_size=10), + callback_routing_key="callback.questions", + session_notifier=session_notifier, + ) + body = _envelope( + { + "sessionId": 99, + "mode": "TECHNICAL", + "jobCategories": ["BACKEND"], + "documents": [], + "maxQuestions": 5, + "initialQuestionCount": 1, + } + ) + await consumer.handle(_StubMessage(body)) + + calls = session_notifier.emit_progress.await_args_list + assert [c.kwargs["phase"] for c in calls] == [ + "CONTEXT_BUILDING", + "GENERATING", + "FINALIZING", + ] + assert all(c.kwargs["event_type"] == "QUESTION_POOL_PROGRESS" for c in calls) + assert all(c.kwargs["session_id"] == 99 for c in calls) + assert all(c.kwargs["trace_id"] == "t-1" for c in calls) + publisher.publish.assert_awaited_once() + + +@pytest.mark.asyncio +async def test_consumer_skips_finalizing_when_generate_fails(): + """생성 실패 시 FAILED 콜백은 나가되, 실패 직전에 "마무리" 진행 문구가 스치면 + 오해를 부르므로 FINALIZING 은 발행하지 않는다.""" + generator = MagicMock() + generator.generate = AsyncMock(side_effect=RuntimeError("llm down")) + publisher = MagicMock() + publisher.publish = AsyncMock() + session_notifier = MagicMock() + session_notifier.emit_progress = AsyncMock() + + consumer = QuestionsConsumer( + generator=generator, + publisher=publisher, + idempotency=LruIdempotencyStore(max_size=10), + callback_routing_key="callback.questions", + session_notifier=session_notifier, + ) + body = _envelope( + { + "sessionId": 99, + "mode": "TECHNICAL", + "jobCategories": ["BACKEND"], + "documents": [], + "maxQuestions": 5, + "initialQuestionCount": 1, + } + ) + await consumer.handle(_StubMessage(body)) + + phases = [c.kwargs["phase"] for c in session_notifier.emit_progress.await_args_list] + assert phases == ["CONTEXT_BUILDING", "GENERATING"] + payload: QuestionPoolCallbackPayload = publisher.publish.await_args.kwargs[ + "payload" + ] + assert payload.status == "FAILED" diff --git a/ai/tests/test_session_notify.py b/ai/tests/test_session_notify.py index 02af3932..57b70d73 100644 --- a/ai/tests/test_session_notify.py +++ b/ai/tests/test_session_notify.py @@ -59,3 +59,67 @@ async def publish(self, **kwargs): await n.emit_audio( session_id=1, message_id=2, seq=0, ext="mp3", duration_sec=None, trace_id="t" ) + + +@pytest.mark.asyncio +async def test_emit_progress_publishes_session_progress_event(): + pub = _FakePublisher() + n = SessionRealtimeNotifier(publisher=pub, routing_key="realtime.session.notify") + await n.emit_progress( + event_type="FEEDBACK_PROGRESS", + session_id=7, + phase="SCORING", + message="세부 평가를 진행하고 있어요. (2/5)", + trace_id="t1", + completed=2, + total=5, + ) + assert len(pub.calls) == 1 + call = pub.calls[0] + assert call["routing_key"] == "realtime.session.notify" + assert call["message_type"] == "realtime.session.notify" + assert call["context"].session_id == 7 + assert call["payload"].event_type == "FEEDBACK_PROGRESS" + data = call["payload"].data + assert (data.session_id, data.phase, data.completed, data.total) == ( + 7, + "SCORING", + 2, + 5, + ) + # RealTime Envelope.payload 직렬화(camelCase) 확인. + dumped = call["payload"].model_dump(by_alias=True) + assert dumped["data"]["sessionId"] == 7 + assert dumped["eventType"] == "FEEDBACK_PROGRESS" + + +@pytest.mark.asyncio +async def test_emit_progress_without_counter_omits_completed_total(): + pub = _FakePublisher() + n = SessionRealtimeNotifier(publisher=pub, routing_key="realtime.session.notify") + await n.emit_progress( + event_type="QUESTION_POOL_PROGRESS", + session_id=9, + phase="GENERATING", + message="첫 질문을 만들고 있어요.", + trace_id="t1", + ) + data = pub.calls[0]["payload"].data + assert data.completed is None + assert data.total is None + + +@pytest.mark.asyncio +async def test_emit_progress_swallows_publish_errors(): + class Boom: + async def publish(self, **kwargs): + raise RuntimeError("down") + + n = SessionRealtimeNotifier(publisher=Boom(), routing_key="realtime.session.notify") + await n.emit_progress( + event_type="QUESTION_POOL_PROGRESS", + session_id=1, + phase="GENERATING", + message="x", + trace_id="t", + ) From 5b2e5b407da9622f6d9cffe87a55f65f5288db92 Mon Sep 17 00:00:00 2001 From: Jaeho Date: Sat, 22 Aug 2026 22:30:28 +0900 Subject: [PATCH 2/3] =?UTF-8?q?feat(frontend):=20=EC=A7=88=EB=AC=B8=20?= =?UTF-8?q?=ED=92=80=C2=B7=ED=94=BC=EB=93=9C=EB=B0=B1=20=EC=83=9D=EC=84=B1?= =?UTF-8?q?=20=EC=A7=84=ED=96=89=20=EB=AC=B8=EA=B5=AC=20=ED=91=9C=EC=8B=9C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 질문 풀: 라이브 WS QUESTION_POOL_PROGRESS → InterviewPreparing 대기 화면 문구 대체 (interviewEvent pool-progress 액션 + useLiveInterview state) - 피드백: 세션 SSE FEEDBACK_PROGRESS → FeedbackReportSkeleton 캡션 (useFeedbackLive progress state, 카운터 completed/total 포함) - 이벤트 미수신 시 기존 기본 안내 문구 폴백 — 유실이 UX 를 깨지 않음 - 워크스페이스 useAnalysisProgress TTL 스토어(user 채널 전용)는 미사용 - 테스트 7건 추가 (이벤트 매핑·진행 state·비정상 페이로드 무시·문구 대체/폴백) --- .../feedback/model/useFeedbackLive.test.tsx | 41 +++++++++++++++++++ .../feedback/model/useFeedbackLive.ts | 30 ++++++++++++-- .../ui/FeedbackReportSkeleton.test.tsx | 14 +++++++ .../feedback/ui/FeedbackReportSkeleton.tsx | 8 ++-- .../interview/model/interviewEvent.test.ts | 3 ++ .../interview/model/interviewEvent.ts | 4 ++ .../interview/model/useLiveInterview.ts | 6 +++ .../ui/live/InterviewPreparing.test.tsx | 26 ++++++++++++ .../interview/ui/live/InterviewPreparing.tsx | 11 ++++- .../interview/ui/live/LiveInterview.tsx | 3 +- .../ui/SessionFeedbackPage.tsx | 7 +++- 11 files changed, 141 insertions(+), 12 deletions(-) create mode 100644 frontend/src/features/interview/ui/live/InterviewPreparing.test.tsx diff --git a/frontend/src/features/feedback/model/useFeedbackLive.test.tsx b/frontend/src/features/feedback/model/useFeedbackLive.test.tsx index 5e92d32b..f87d2e3d 100644 --- a/frontend/src/features/feedback/model/useFeedbackLive.test.tsx +++ b/frontend/src/features/feedback/model/useFeedbackLive.test.tsx @@ -93,4 +93,45 @@ describe('useFeedbackLive', () => { expect(result.current.data?.overallScore).toBe(80) expect(vi.mocked(getFeedback).mock.calls.length).toBe(callsAfterLoad) }) + it('FEEDBACK_PROGRESS 수신 시 진행 문구·카운터를 노출한다 (B2)', async () => { + vi.mocked(getFeedback).mockRejectedValue(notReady()) + + const { result } = setup() + await waitFor(() => expect(FakeES.last).not.toBeNull()) + expect(result.current.progress).toBeNull() + + act(() => + FakeES.last?.emit('FEEDBACK_PROGRESS', { + data: { sessionId: 99, phase: 'SCORING', message: '세부 평가를 진행하고 있어요. (2/5)', completed: 2, total: 5 }, + }), + ) + expect(result.current.progress).toEqual({ + message: '세부 평가를 진행하고 있어요. (2/5)', + completed: 2, + total: 5, + }) + + // 카운터 없는 진행 이벤트(PREPARING/FINALIZING)는 message 만 채워진다. + act(() => + FakeES.last?.emit('FEEDBACK_PROGRESS', { + data: { sessionId: 99, phase: 'FINALIZING', message: '피드백 리포트를 정리하고 있어요.' }, + }), + ) + expect(result.current.progress).toEqual({ + message: '피드백 리포트를 정리하고 있어요.', + completed: null, + total: null, + }) + }) + + it('message 없는 비정상 FEEDBACK_PROGRESS 페이로드는 무시한다', async () => { + vi.mocked(getFeedback).mockRejectedValue(notReady()) + + const { result } = setup() + await waitFor(() => expect(FakeES.last).not.toBeNull()) + + act(() => FakeES.last?.emit('FEEDBACK_PROGRESS', { data: { completed: 1 } })) + act(() => FakeES.last?.emit('FEEDBACK_PROGRESS', {})) + expect(result.current.progress).toBeNull() + }) }) diff --git a/frontend/src/features/feedback/model/useFeedbackLive.ts b/frontend/src/features/feedback/model/useFeedbackLive.ts index 8ee03ec1..ab0a3cce 100644 --- a/frontend/src/features/feedback/model/useFeedbackLive.ts +++ b/frontend/src/features/feedback/model/useFeedbackLive.ts @@ -1,11 +1,21 @@ -import { useCallback, useEffect, useRef } from 'react' +import { useCallback, useEffect, useRef, useState } from 'react' import { useQuery, useQueryClient } from '@tanstack/react-query' import { useEventStream } from '@/shared/hooks' import type { StreamConnectionStatus } from '@/shared/hooks' import { getFeedback } from '../api/feedbackApi' import { feedbackKeys, isFeedbackPending } from './useFeedback' -// 피드백 조회 + 세션 채널 SSE(FEEDBACK_READY) 구독. +// AI 피드백 생성 진행 문구(FEEDBACK_PROGRESS, 휘발성) — 스켈레톤 캡션에 표시. +export type FeedbackProgress = { + message: string + completed: number | null + total: number | null +} + +// RealTime SSE data 봉투: { data: , traceId } (useWorkspaceAnalysisStream 과 동일). +type StreamData = { data?: T; traceId?: string | null } + +// 피드백 조회 + 세션 채널 SSE(FEEDBACK_READY·FEEDBACK_PROGRESS) 구독. // - 준비 완료 이벤트 수신 즉시 재조회 — 3s 폴링 간격을 기다리지 않는다. // - 기존 3s×40 재시도는 백스톱으로 축소: 스트림이 살아 있으면 15s 간격(이벤트가 주 신호), // 죽어 있으면 기존대로 3s (api-conventions §9 — SSE 우선, 끊기면 폴링). @@ -34,16 +44,28 @@ export function useFeedbackLive(sessionId: number, getToken: () => Promise(null) + const onProgress = useCallback((raw: unknown) => { + const data = (raw as StreamData<{ message?: string; completed?: number; total?: number }> | null) + ?.data + if (!data || typeof data.message !== 'string') return + setProgress({ + message: data.message, + completed: typeof data.completed === 'number' ? data.completed : null, + total: typeof data.total === 'number' ? data.total : null, + }) + }, []) + // 피드백이 도착하면 스트림을 닫는다 (path: null). const streamStatus = useEventStream({ path: query.data ? null : `/realtime/stream/sessions/${sessionId}`, getToken, - handlers: { FEEDBACK_READY: onReady }, + handlers: { FEEDBACK_READY: onReady, FEEDBACK_PROGRESS: onProgress }, }) useEffect(() => { statusRef.current = streamStatus }, [streamStatus]) - return { ...query, streamStatus } + return { ...query, streamStatus, progress } } diff --git a/frontend/src/features/feedback/ui/FeedbackReportSkeleton.test.tsx b/frontend/src/features/feedback/ui/FeedbackReportSkeleton.test.tsx index 3f7f42e7..6a086b15 100644 --- a/frontend/src/features/feedback/ui/FeedbackReportSkeleton.test.tsx +++ b/frontend/src/features/feedback/ui/FeedbackReportSkeleton.test.tsx @@ -22,4 +22,18 @@ describe('FeedbackReportSkeleton', () => { }) expect(screen.getByRole('status')).toHaveTextContent('(12s 경과)') }) + it('진행 이벤트가 오면 기본 안내 대신 진행 문구를 보여준다 (B2)', () => { + render( + , + ) + expect(screen.getByRole('status')).toHaveTextContent('세부 평가를 진행하고 있어요. (3/5)') + expect(screen.getByRole('status')).not.toHaveTextContent('보통 1분 내외') + }) + + it('진행 이벤트 미수신이면 기본 안내 문구를 유지한다', () => { + render() + expect(screen.getByRole('status')).toHaveTextContent('보통 1분 내외') + }) }) diff --git a/frontend/src/features/feedback/ui/FeedbackReportSkeleton.tsx b/frontend/src/features/feedback/ui/FeedbackReportSkeleton.tsx index a6b007d4..ed762d49 100644 --- a/frontend/src/features/feedback/ui/FeedbackReportSkeleton.tsx +++ b/frontend/src/features/feedback/ui/FeedbackReportSkeleton.tsx @@ -1,12 +1,14 @@ import { useEffect, useState } from 'react' +import type { FeedbackProgress } from '../model/useFeedbackLive' function Pulse({ className }: { className: string }) { return
} // 피드백 생성 대기 화면 — 리포트가 올 자리의 형태를 미리 보여주고 경과 시간을 알린다. -// 정적 텍스트만 있으면 "멈췄나?" 로 읽힌다 (A4). 완료는 FEEDBACK_READY SSE 로 즉시 반영된다. -export function FeedbackReportSkeleton() { +// 정적 텍스트만 있으면 "멈췄나?" 로 읽힌다 (A4). 완료는 FEEDBACK_READY SSE 로 즉시 반영되고, +// 생성 단계는 FEEDBACK_PROGRESS SSE 진행 문구로 알린다 (B2). 이벤트 미수신이면 기본 문구. +export function FeedbackReportSkeleton({ progress }: { progress?: FeedbackProgress | null }) { const [elapsedSec, setElapsedSec] = useState(0) useEffect(() => { const id = setInterval(() => setElapsedSec((s) => s + 1), 1_000) @@ -18,7 +20,7 @@ export function FeedbackReportSkeleton() {

피드백을 생성하는 중입니다…

- 답변을 종합 분석하고 있어요. 보통 1분 내외, 길면 2분 이상 걸릴 수 있습니다. + {progress?.message ?? '답변을 종합 분석하고 있어요. 보통 1분 내외, 길면 2분 이상 걸릴 수 있습니다.'}{' '} 완성되면 자동으로 표시됩니다. ({elapsedSec}s 경과)

diff --git a/frontend/src/features/interview/model/interviewEvent.test.ts b/frontend/src/features/interview/model/interviewEvent.test.ts index 1b12768b..02007da3 100644 --- a/frontend/src/features/interview/model/interviewEvent.test.ts +++ b/frontend/src/features/interview/model/interviewEvent.test.ts @@ -23,4 +23,7 @@ describe('interviewEventAction', () => { it('SESSION_MESSAGE_AUDIO → 오디오 큐', () => { expect(interviewEventAction('SESSION_MESSAGE_AUDIO')).toEqual({ kind: 'queue-audio' }) }) + it('QUESTION_POOL_PROGRESS → 질문 풀 진행 문구 (B2)', () => { + expect(interviewEventAction('QUESTION_POOL_PROGRESS')).toEqual({ kind: 'pool-progress' }) + }) }) diff --git a/frontend/src/features/interview/model/interviewEvent.ts b/frontend/src/features/interview/model/interviewEvent.ts index 367198e0..ca064659 100644 --- a/frontend/src/features/interview/model/interviewEvent.ts +++ b/frontend/src/features/interview/model/interviewEvent.ts @@ -4,6 +4,7 @@ export type InterviewAction = | { kind: 'redirect-feedback' } | { kind: 'append-delta' } | { kind: 'queue-audio' } + | { kind: 'pool-progress' } | { kind: 'rollback-optimistic' } | { kind: 'ignore' } @@ -19,6 +20,9 @@ export function interviewEventAction(event: string): InterviewAction { return { kind: 'append-delta' } case 'SESSION_MESSAGE_AUDIO': return { kind: 'queue-audio' } + // AI 가 질문 풀 생성 단계를 세션 채널로 직접 발행(휘발성) — 대기 화면 진행 문구용. + case 'QUESTION_POOL_PROGRESS': + return { kind: 'pool-progress' } // realtime 서버는 제출 실패 시 'error' 프레임을 보낸다(ws.go outboundFrame). // 백엔드 SseEventType.ERROR 도 동일 취급. 무시하면 거부된 답변이 영구 '전송 중'으로 남는다. case 'error': diff --git a/frontend/src/features/interview/model/useLiveInterview.ts b/frontend/src/features/interview/model/useLiveInterview.ts index b088dadb..581b4b9f 100644 --- a/frontend/src/features/interview/model/useLiveInterview.ts +++ b/frontend/src/features/interview/model/useLiveInterview.ts @@ -39,6 +39,8 @@ export function useLiveInterview(sessionId: number, deliveryMode: DeliveryMode = // 전송 실패로 롤백된 답변 본문 — 컴포저가 입력창을 복원하는 데 사용(nonce 로 매 실패마다 트리거). const [restoreDraft, setRestoreDraft] = useState<{ content: string; nonce: number } | null>(null) const [deltaBuffer, setDeltaBuffer] = useState({}) + // 질문 풀 생성 진행 문구(QUESTION_POOL_PROGRESS, 휘발성) — 대기 화면(InterviewPreparing)에서만 사용. + const [poolProgress, setPoolProgress] = useState(null) // 라이브 세그먼트 오디오가 지금 재생 중인 메시지(아바타·질문 카드의 '말하는 중' 표시용). const [speakingAudio, setSpeakingAudio] = useState<{ msgId: number | null; playing: boolean }>({ msgId: null, @@ -217,6 +219,9 @@ export function useLiveInterview(sessionId: number, deliveryMode: DeliveryMode = if (payload && typeof payload.messageId === 'number' && typeof payload.seq === 'number') { setDeltaBuffer((b) => applyDelta(b, payload)) } + } else if (action.kind === 'pool-progress') { + const p = (frame.data as { data?: { message?: string } } | undefined)?.data + if (p && typeof p.message === 'string') setPoolProgress(p.message) } else if (action.kind === 'queue-audio') { // 텍스트 모드에서는 음성을 자동재생하지 않는다(읽는 도중 끼어들기 방지). // 수동 재생은 StageQuestion 의 전체 파일(callback.tts) 경로로 제공된다. @@ -349,5 +354,6 @@ export function useLiveInterview(sessionId: number, deliveryMode: DeliveryMode = wasSegmented, isSpeaking, firstQuestionReady, + poolProgress, } } diff --git a/frontend/src/features/interview/ui/live/InterviewPreparing.test.tsx b/frontend/src/features/interview/ui/live/InterviewPreparing.test.tsx new file mode 100644 index 00000000..bba0c01f --- /dev/null +++ b/frontend/src/features/interview/ui/live/InterviewPreparing.test.tsx @@ -0,0 +1,26 @@ +import { describe, it, expect } from 'vitest' +import { render, screen } from '@testing-library/react' +import type { Session } from '@/domain/session' +import { InterviewPreparing } from './InterviewPreparing' + +const session = { + id: 1, + title: '백엔드 모의 면접', + maxQuestions: 5, + totalQuestionCount: 0, +} as Session + +describe('InterviewPreparing', () => { + it('진행 이벤트 미수신이면 기본 안내 문구를 보여준다', () => { + render() + expect(screen.getByRole('status')).toHaveTextContent('첫 질문을 만들고 있어요') + }) + + it('QUESTION_POOL_PROGRESS 진행 문구가 오면 기본 문구를 대체한다 (B2)', () => { + render( + , + ) + expect(screen.getByRole('status')).toHaveTextContent('자료를 바탕으로 첫 질문을 만들고 있어요.') + expect(screen.getByRole('status')).not.toHaveTextContent('준비가 끝나면 바로 면접이 시작됩니다') + }) +}) diff --git a/frontend/src/features/interview/ui/live/InterviewPreparing.tsx b/frontend/src/features/interview/ui/live/InterviewPreparing.tsx index 3baa52d1..6c4a16fb 100644 --- a/frontend/src/features/interview/ui/live/InterviewPreparing.tsx +++ b/frontend/src/features/interview/ui/live/InterviewPreparing.tsx @@ -4,7 +4,14 @@ import { Heading } from '@/shared/ui' // 면접 시작 직후, 첫 질문이 준비될 때까지 스테이지 진입 전에 머무는 대기 화면. // 첫 질문이 도착하면 LiveInterview 가 InterviewStage 로 전환하며 화면이 "켜진다". -export function InterviewPreparing({ session }: { session: Session }) { +export function InterviewPreparing({ + session, + progressMessage, +}: { + session: Session + /** AI 질문 풀 생성 진행 문구(QUESTION_POOL_PROGRESS). 없으면 기본 안내 문구. */ + progressMessage?: string | null +}) { const progress = sessionProgress(session) return (
@@ -21,7 +28,7 @@ export function InterviewPreparing({ session }: { session: Session }) { {session.title ?? '모의 면접'}

- 첫 질문을 만들고 있어요. 준비가 끝나면 바로 면접이 시작됩니다. + {progressMessage ?? '첫 질문을 만들고 있어요. 준비가 끝나면 바로 면접이 시작됩니다.'}

총 {progress.max}개의 질문이 준비됩니다.

diff --git a/frontend/src/features/interview/ui/live/LiveInterview.tsx b/frontend/src/features/interview/ui/live/LiveInterview.tsx index 6a264892..e2c15601 100644 --- a/frontend/src/features/interview/ui/live/LiveInterview.tsx +++ b/frontend/src/features/interview/ui/live/LiveInterview.tsx @@ -28,6 +28,7 @@ export function LiveInterview({ sessionId }: { sessionId: number }) { wasSegmented, isSpeaking, firstQuestionReady, + poolProgress, } = useLiveInterview(sessionId, deliveryMode) // 에러 분기가 스피너 분기보다 먼저 와야 한다 — 쿼리 실패 시 isLoading 은 false 이고 @@ -70,7 +71,7 @@ export function LiveInterview({ sessionId }: { sessionId: number }) { } // 면접은 시작됐지만 첫 질문이 아직 안 왔으면 스테이지 진입 전 대기 화면을 보여준다. if (!firstQuestionReady) { - return + return } const awaitingQuestion = turn === 'WAITING_FOR_QUESTION' diff --git a/frontend/src/pages/SessionFeedback/ui/SessionFeedbackPage.tsx b/frontend/src/pages/SessionFeedback/ui/SessionFeedbackPage.tsx index 5eb1a0f6..2ce29d69 100644 --- a/frontend/src/pages/SessionFeedback/ui/SessionFeedbackPage.tsx +++ b/frontend/src/pages/SessionFeedback/ui/SessionFeedbackPage.tsx @@ -24,7 +24,10 @@ export default function SessionFeedbackPage() { // 피드백 조회 + FEEDBACK_READY SSE — 준비 완료 즉시 표시(폴링은 백스톱). // 세션 stream token 은 interview 슬라이스 소유라 페이지가 주입한다. const getToken = useCallback(() => fetchSessionStreamToken(sessionId), [sessionId]) - const { data, isLoading, isError, error, refetch } = useFeedbackLive(sessionId, getToken) + const { data, isLoading, isError, error, refetch, progress } = useFeedbackLive( + sessionId, + getToken, + ) const regenerate = useRegenerateFeedback(sessionId) // 재도전은 원본 세션의 자료 수를 알아야 "몇 개가 빠졌는지" 안내할 수 있다. const { data: session } = useSession(sessionId) @@ -61,7 +64,7 @@ export default function SessionFeedbackPage() { } /> - {isLoading && } + {isLoading && } {isError && (isFeedbackPending(error) ? ( From 5c8e09c01fb41c7574078c5ec1a6f6c2f62d5d60 Mon Sep 17 00:00:00 2001 From: Jaeho Date: Sat, 22 Aug 2026 22:30:28 +0900 Subject: [PATCH 3/3] =?UTF-8?q?docs:=20=EC=A7=88=EB=AC=B8=20=ED=92=80?= =?UTF-8?q?=C2=B7=ED=94=BC=EB=93=9C=EB=B0=B1=20=EC=83=9D=EC=84=B1=20?= =?UTF-8?q?=EC=A7=84=ED=96=89=20=EC=9D=B4=EB=B2=A4=ED=8A=B8=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.12-3 신설 — realtime.session.notify 발행 스펙(phase·카운터) - event-stream.md §2 이벤트 목록 + §3.3-3 신설 (payload 예시·소비처·폴백) - ai/·frontend/ CLAUDE.md 현재 상태 갱신 --- ai/CLAUDE.md | 9 +++++++++ docs/event-stream.md | 15 ++++++++++++++- docs/messaging.md | 8 ++++++++ frontend/CLAUDE.md | 1 + 4 files changed, 32 insertions(+), 1 deletion(-) diff --git a/ai/CLAUDE.md b/ai/CLAUDE.md index 95702976..fcdc88a4 100644 --- a/ai/CLAUDE.md +++ b/ai/CLAUDE.md @@ -397,4 +397,13 @@ docker run --env-file .env -p 8000:8000 stackup-ai 하드 타임아웃을 추가하고, 모든 `ChatOpenAI` 호출에 `llm_pro_timeout_sec`(30s)/`llm_flash_timeout_sec` (10s) 요청 타임아웃을 명시했다(이전엔 미설정 — SDK 기본값까지 무기한 대기 가능). +- **질문 풀·피드백 생성 진행 이벤트 본 구현 (B2)**: 스트리밍이 없는 두 블로킹 생성 경로(질문 풀 Pro ≤30s, + 피드백 병렬 gather ≈2분 예산)가 진행 중 무통보였던 것을 고쳤다. `SessionRealtimeNotifier.emit_progress` + (`messaging/session_notify.py`)가 `realtime.session.notify` 로 `QUESTION_POOL_PROGRESS`/`FEEDBACK_PROGRESS` + 를 직접 발행(휘발성, 실패는 경고만). `questions_consumer` 는 CONTEXT_BUILDING→GENERATING→FINALIZING 순차 + 3단계, `feedback_consumer` 는 최상위가 `asyncio.gather` 병렬이라 순차 phase 대신 **태스크 완료 카운터** + (SCORING `completed/total`, 각 세부 평가 완료 시 emit)로 표현한다. 기존 `AnalysisProgressNotifier` + (user 채널)는 무변경 — user 채널 경로 불변은 회귀 테스트(`tests/test_progress.py`)로 고정. + 와이어링은 `runner.py` 에서 followup 과 동일 notifier 인스턴스 재사용. 스펙: `docs/messaging.md §5.12-3`. + 각 도입 시 본 문서 갱신. diff --git a/docs/event-stream.md b/docs/event-stream.md index f2293b8b..79ecc007 100644 --- a/docs/event-stream.md +++ b/docs/event-stream.md @@ -40,7 +40,7 @@ WS /realtime/sessions/{sessionId}/audio # RT3 실시간 음성 답변 스트 ## 2. 이벤트 포맷 -> - `event` 이름은 **`SseEventType` enum 이름(대문자)**: `DOC_STATE`·`REPO_STATE`·`SESSION_MESSAGE`·`SESSION_STATE`·`FEEDBACK_READY`·`ERROR`·`KEEP_ALIVE` (+ AI 가 직접 발행하는 `ANALYSIS_PROGRESS`·`SESSION_MESSAGE_DELTA`·`SESSION_MESSAGE_AUDIO`). **`session.message` 같은 소문자 점표기가 아니다.** (RealTime `bridge/dispatcher.go` `Type: env.Payload.EventType` → `sse.go`/`ws.go` 가 그대로 전달.) RealTime 디스패처는 이벤트 타입을 화이트리스트 없이 투명 전달하므로, AI 가 새 이벤트 타입(`SESSION_MESSAGE_DELTA` 등)을 발행해도 RealTime 코드 변경이 필요 없다. +> - `event` 이름은 **`SseEventType` enum 이름(대문자)**: `DOC_STATE`·`REPO_STATE`·`SESSION_MESSAGE`·`SESSION_STATE`·`FEEDBACK_READY`·`ERROR`·`KEEP_ALIVE` (+ AI 가 직접 발행하는 `ANALYSIS_PROGRESS`·`SESSION_MESSAGE_DELTA`·`SESSION_MESSAGE_AUDIO`·`QUESTION_POOL_PROGRESS`·`FEEDBACK_PROGRESS`). **`session.message` 같은 소문자 점표기가 아니다.** (RealTime `bridge/dispatcher.go` `Type: env.Payload.EventType` → `sse.go`/`ws.go` 가 그대로 전달.) RealTime 디스패처는 이벤트 타입을 화이트리스트 없이 투명 전달하므로, AI 가 새 이벤트 타입(`SESSION_MESSAGE_DELTA` 등)을 발행해도 RealTime 코드 변경이 필요 없다. > - `data` 봉투는 `{"data": , "traceId": "..."}` 다 (`realtime/CLAUDE.md §8`). payload 필드는 camelCase. > - 클라는 SSE `addEventListener(, …)` / WS `frame.event === ''` 로 매칭한다. @@ -141,6 +141,19 @@ WS(RT1)는 같은 내용을 JSON 한 줄 프레임으로: `{ "id": , "e - 세그먼트는 S3 `interview/tts/{sessionId}/{messageId}/seg-{seq}.{ext}` 에 저장(DB 미기록, 휘발성). 프론트는 Core 프록시 `GET /api/sessions/{sid}/messages/{mid}/audio/segments/{seq}?ext=` 로 받아 **seq 순서대로 순차 재생**. - `DONT_KNOW` 면 발행 안 함. 라이브 세그먼트를 재생한 메시지는 완료 후 whole-message TTS **autoPlay 억제**(중복 재생 방지) — 수동 "다시 듣기"만 동작. +### 3.3-3 질문 풀·피드백 생성 진행 (`event: QUESTION_POOL_PROGRESS` / `FEEDBACK_PROGRESS`) + +**휘발성** 이벤트. 질문 풀 생성(Pro, 최대 30s)·피드백 생성(≈2분 예산) 대기 화면에 진행 문구를 흘린다. **AI 서버가 `stackup.realtime`(`realtime.session.notify`)으로 직접 발행**(Core·DB 미경유) → 세션 채널(SSE·WS 공통)로 fan-out. 발행 스펙: [`messaging.md §5.12-3`](./messaging.md). + +```json +{ "data": { "sessionId": 99, "phase": "GENERATING", "message": "자료를 바탕으로 첫 질문을 만들고 있어요." }, "traceId": "..." } +{ "data": { "sessionId": 99, "phase": "SCORING", "message": "세부 평가를 진행하고 있어요. (2/5)", "completed": 2, "total": 5 }, "traceId": "..." } +``` + +- `QUESTION_POOL_PROGRESS` phase: `CONTEXT_BUILDING → GENERATING → FINALIZING` (순차). 소비: 라이브 면접 WS(`interviewEvent.ts` → `InterviewPreparing` 대기 화면 문구). +- `FEEDBACK_PROGRESS` phase: `PREPARING → SCORING(완료 카운터 completed/total, 병렬 평가라 순차 아님) → FINALIZING`. 소비: 피드백 페이지 세션 SSE(`useFeedbackLive` → `FeedbackReportSkeleton` 캡션). +- 프론트 처리: 쿼리 무효화 없이 로컬 state 만 갱신(`message` 그대로 표시). 이벤트가 하나도 안 와도 기본 안내 문구가 유지되므로 유실은 UX 저하 없이 흡수된다. 종료 신호는 기존 `SESSION_MESSAGE`/`FEEDBACK_READY` 가 정본. + ### 3.4 세션 상태 (`event: SESSION_STATE`) ```json { diff --git a/docs/messaging.md b/docs/messaging.md index 9c162b50..ec3c54d0 100644 --- a/docs/messaging.md +++ b/docs/messaging.md @@ -551,6 +551,14 @@ AI followup consumer 가 토큰 스트림 중 문장 경계마다 그 문장만 - 세그먼트 S3 키 규칙(AI·Core 공유): `interview/tts/{sessionId}/{messageId}/seg-{seq}.{ext}`. Core 는 DB 미기록, `GET …/messages/{mid}/audio/segments/{seq}?ext=` 프록시에서 소유권 검증 후 규칙으로 키 재구성(ext 화이트리스트 wav|mp3|ogg|m4a). - 상세 SSE 스펙: [`event-stream.md §3.2-1`](./event-stream.md). +### 5.12-3 `realtime.session.notify` — 질문 풀·피드백 생성 진행 (AI 직접 발행) +질문 풀 생성(Pro, 최대 30s)과 피드백 생성(패널 병렬, ≈2분 예산)은 진행 중 무통보 블로킹이었다. `ANALYSIS_PROGRESS` 와 동일 패턴으로 AI 서버가 생성 단계를 `realtime.session.notify` 로 **직접** 발행(Core·DB 미경유, 휘발성 — 발행 실패는 경고만)해 세션 채널로 흘린다. +- 발행: AI `SessionRealtimeNotifier.emit_progress` (`questions_consumer`/`feedback_consumer` 에서 호출). `context.sessionId` 로 세션 채널 라우팅. +- `payload.eventType` = `QUESTION_POOL_PROGRESS`, `payload.data` = `{ sessionId, phase, message }`. phase ∈ `CONTEXT_BUILDING | GENERATING | FINALIZING` (순차). +- `payload.eventType` = `FEEDBACK_PROGRESS`, `payload.data` = `{ sessionId, phase, message, completed?, total? }`. phase ∈ `PREPARING | SCORING | FINALIZING`. 세부 평가 5개가 `asyncio.gather` 병렬이라 SCORING 은 순차 단계가 아닌 **완료 카운터**(`completed`/`total`, 시작 시 0) 로 표현한다. +- 종료 정본은 기존대로 Core 콜백(`callback.questions`/`callback.feedback`) → `SESSION_MESSAGE`/`FEEDBACK_READY`. 진행 이벤트 유실은 UI 기본 문구 폴백으로 흡수(프론트 `InterviewPreparing`/`FeedbackReportSkeleton`). +- 상세 SSE 스펙: [`event-stream.md §3.3-3`](./event-stream.md). + ### 5.14 발행 시점 규약 — 반드시 커밋 후에 발행한다 **Core 의 모든 작업 요청 발행(`stackup.core-to-ai`)은 DB 커밋 이후에 일어나야 한다.** diff --git a/frontend/CLAUDE.md b/frontend/CLAUDE.md index 6853042e..50952d27 100644 --- a/frontend/CLAUDE.md +++ b/frontend/CLAUDE.md @@ -237,6 +237,7 @@ SEED 팔레트 블록에는 `prefers-color-scheme` 미디어쿼리가 없어서, - **SSE + WebSocket 병행** — 작업 상태 푸시(분석·피드백)는 SSE, 라이브 면접 메시지는 WS(`features/interview/model/useInterviewSocket.ts`). (루트 CLAUDE.md §8 과 동일) - 구현: `shared/hooks/useEventStream.ts` — 자동 재연결(지수 백오프) + 연결 상태 반환. 워크스페이스는 단절(closed) 시 배너 표시 + 목록 쿼리 5s 폴백 폴링(`useAnalysisFallbackPolling`) +- 생성 진행 문구(휘발성, B2): 질문 풀 대기 화면은 WS `QUESTION_POOL_PROGRESS`(`interviewEvent.ts` → `InterviewPreparing`), 피드백 대기 스켈레톤은 세션 SSE `FEEDBACK_PROGRESS`(`useFeedbackLive` → `FeedbackReportSkeleton`). 둘 다 로컬 state 로만 표시하고 이벤트 미수신 시 기본 안내 문구로 폴백 — 워크스페이스의 `useAnalysisProgress` TTL 스토어(90s, user 채널 전용)는 통과하지 않는다 - 미디어 스트림(음성/영상)만 WebRTC: `features/interview/lib/media/` - 이벤트 스펙: [`/docs/event-stream.md`](../docs/event-stream.md)