Skip to content

Commit be4c5aa

Browse files
authored
Merge pull request #206 from Team-StackUp/feat/feedback-failed-callback
feat: 피드백 생성 실패 신호 — FAILED 콜백으로 무기한 대기 제거
2 parents 5542169 + 9a69374 commit be4c5aa

11 files changed

Lines changed: 456 additions & 103 deletions

File tree

ai/CLAUDE.md

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -397,6 +397,20 @@ docker run --env-file .env -p 8000:8000 stackup-ai
397397
하드 타임아웃을 추가하고, 모든 `ChatOpenAI` 호출에 `llm_pro_timeout_sec`(30s)/`llm_flash_timeout_sec`
398398
(10s) 요청 타임아웃을 명시했다(이전엔 미설정 — SDK 기본값까지 무기한 대기 가능).
399399

400+
- **피드백 생성 실패 신호 본 구현**: `feedback_consumer` 만 위 리팩터에서 빠져 있던 gap 을 닫았다 —
401+
패널·부가 평가는 내부 폴백(빈 결과/생략)으로 흡수되지만, 그 방어망 밖(트랜스크립트/RAG 컨텍스트 빌드,
402+
payload 조립, 발행 등)의 예상 못 한 예외는 그대로 새서 DLQ 로만 격리되고 Core 는 아무 신호도 못 받아
403+
세션이 "피드백 생성 중"에 무기한 멈췄다. `handle()` 본문을 payload 를 **반환**하는 `_process()`
404+
추출하고 envelope 파싱·멱등 체크 이후의 생성 전 구간을 try/except 로 감싸, 실패 시
405+
`FeedbackCallbackPayload`(`status=FAILED`, `errorCode`, `errorMessage`(상한 500자), `retriable`)를
406+
발행하고 ack 한다(`_publish_failed`). 분류는 questions/followup 과 동일 — `TypeError`
407+
`GENERATION_SCHEMA_INVALID`/`retriable=false`, 그 외 `UNEXPECTED`/`retriable=true`. **성공 콜백
408+
발행 실패는 생성 실패가 아니다** — FAILED 오인 발행 없이 원 예외로 DLQ(변경 전과 동일, 재처리
409+
가능). 콜백을 하나도 못 낸 채 DLQ 로 가는 경로(폴백 발행 실패 포함)는 `LruIdempotencyStore.unmark`
410+
로 마킹을 되돌려 재주입 시 duplicate skip 으로 삼켜지지 않게 한다. `status``GenerationStatus`
411+
Literal 재사용, 기본값 `OK` 라 성공 콜백·구버전 소비자와 하위호환.
412+
Core 쪽 처리는 [`backend/CLAUDE.md`](../backend/CLAUDE.md) 참고.
413+
400414
- **질문 풀·피드백 생성 진행 이벤트 본 구현 (B2)**: 스트리밍이 없는 두 블로킹 생성 경로(질문 풀 Pro ≤30s,
401415
피드백 병렬 gather ≈2분 예산)가 진행 중 무통보였던 것을 고쳤다. `SessionRealtimeNotifier.emit_progress`
402416
(`messaging/session_notify.py`)가 `realtime.session.notify``QUESTION_POOL_PROGRESS`/`FEEDBACK_PROGRESS`

ai/src/ai_server/messaging/consumers/feedback_consumer.py

Lines changed: 161 additions & 100 deletions
Original file line numberDiff line numberDiff line change
@@ -124,127 +124,188 @@ async def handle(self, message: AbstractIncomingMessage) -> None:
124124
)
125125
return
126126

127-
req = envelope.payload
127+
try:
128+
payload = await self._process(envelope)
129+
except Exception as exc: # noqa: BLE001
130+
await self._publish_failed(envelope, exc)
131+
return
132+
133+
# 성공 payload 의 발행 실패는 생성 실패가 아니다 — FAILED 콜백으로 오인 발행하지
134+
# 않고 원 예외로 DLQ 에 보내 재처리 가능하게 남긴다(실패 신호 도입 전과 동일 동작).
135+
try:
136+
await self._publish_callback(envelope, payload)
137+
except Exception:
138+
self._idempotency.unmark(envelope.message_id)
139+
raise
128140
log.info(
129-
"feedback.generate.start",
141+
"feedback.generate.done",
130142
message_id=envelope.message_id,
131-
session_id=req.session_id,
132-
msg_count=len(req.messages),
133-
ctx_count=len(req.context_document_ids),
143+
session_id=envelope.payload.session_id,
134144
trace_id=envelope.trace_id,
135145
)
136146

137-
await self._emit_progress(
138-
session_id=req.session_id,
139-
phase="PREPARING",
140-
message="면접 기록을 정리하고 있어요.",
141-
trace_id=envelope.trace_id,
142-
)
143-
transcript = _build_transcript(req.messages)
144-
score_basis = _build_score_basis(req.messages)
145-
rag_context = await self._build_rag_context(req)
146-
voice_analysis_summary = _build_voice_analysis_summary(
147-
req.voice_analysis_summary
148-
)
147+
async def _process(
148+
self, envelope: Envelope[GenerateFeedbackRequest]
149+
) -> FeedbackCallbackPayload:
150+
req = envelope.payload
151+
log.info(
152+
"feedback.generate.start",
153+
message_id=envelope.message_id,
154+
session_id=req.session_id,
155+
msg_count=len(req.messages),
156+
ctx_count=len(req.context_document_ids),
157+
trace_id=envelope.trace_id,
158+
)
149159

150-
# 세부 평가 5개가 병렬(gather)이라 순차 phase 로는 진행을 표현할 수 없다 —
151-
# 각 태스크 완료 시점에 completed/total 카운터로 emit 한다.
152-
scoring_total = 5
153-
scoring_done = 0
160+
await self._emit_progress(
161+
session_id=req.session_id,
162+
phase="PREPARING",
163+
message="면접 기록을 정리하고 있어요.",
164+
trace_id=envelope.trace_id,
165+
)
166+
transcript = _build_transcript(req.messages)
167+
score_basis = _build_score_basis(req.messages)
168+
rag_context = await self._build_rag_context(req)
169+
voice_analysis_summary = _build_voice_analysis_summary(
170+
req.voice_analysis_summary
171+
)
154172

155-
async def _tracked(coro: Awaitable[T]) -> T:
156-
nonlocal scoring_done
157-
task_result = await coro
158-
scoring_done += 1
159-
await self._emit_progress(
160-
session_id=req.session_id,
161-
phase="SCORING",
162-
message=f"세부 평가를 진행하고 있어요. ({scoring_done}/{scoring_total})",
163-
trace_id=envelope.trace_id,
164-
completed=scoring_done,
165-
total=scoring_total,
166-
)
167-
return task_result
173+
# 세부 평가 5개가 병렬(gather)이라 순차 phase 로는 진행을 표현할 수 없다 —
174+
# 각 태스크 완료 시점에 completed/total 카운터로 emit 한다.
175+
scoring_total = 5
176+
scoring_done = 0
168177

178+
async def _tracked(coro: Awaitable[T]) -> T:
179+
nonlocal scoring_done
180+
task_result = await coro
181+
scoring_done += 1
169182
await self._emit_progress(
170183
session_id=req.session_id,
171184
phase="SCORING",
172-
message="평가위원들이 답변을 검토하고 있어요.",
185+
message=f"세부 평가를 진행하고 있어요. ({scoring_done}/{scoring_total})",
173186
trace_id=envelope.trace_id,
174-
completed=0,
187+
completed=scoring_done,
175188
total=scoring_total,
176189
)
190+
return task_result
191+
192+
await self._emit_progress(
193+
session_id=req.session_id,
194+
phase="SCORING",
195+
message="평가위원들이 답변을 검토하고 있어요.",
196+
trace_id=envelope.trace_id,
197+
completed=0,
198+
total=scoring_total,
199+
)
177200

178-
# 종합 피드백 + 자기소개 첫인상 + 직무 적합도(직무 맞춤 모드)를 병렬 실행.
179-
# 첫인상·직무 적합도는 종합 점수(overall)에 미포함 — generator 가 모른 채 계산한 뒤 표시용으로 덧붙인다.
180-
(
181-
result,
182-
self_intro_item,
183-
job_fit_items,
184-
personality_item,
185-
answer_coaching,
186-
) = await asyncio.gather(
187-
_tracked(
188-
self._generate_panel(
189-
job_category=req.job_category,
190-
mode=req.mode,
191-
total_question_count=req.total_question_count,
192-
end_reason=req.end_reason,
193-
transcript=transcript,
194-
score_basis=score_basis,
195-
rag_context=rag_context,
196-
voice_analysis_summary=voice_analysis_summary,
197-
domain_question_counts=req.domain_question_counts,
198-
session_id=req.session_id,
199-
)
200-
),
201-
_tracked(self._evaluate_self_intro(req, voice_analysis_summary)),
202-
_tracked(self._evaluate_job_fit(req, transcript, rag_context)),
203-
_tracked(self._evaluate_personality(req)),
204-
_tracked(self._coach_answers(req)),
205-
)
206-
# 빈 평가위원 항목(점수·내용 모두 없음)은 표시하지 않는다 — LLM 부분 응답이 빈 패널로 새는 것 방지.
207-
extras = [self_intro_item, *job_fit_items, personality_item]
208-
result.panel_breakdown.extend(
209-
e for e in extras if e is not None and _panel_has_content(e)
210-
)
201+
# 종합 피드백 + 자기소개 첫인상 + 직무 적합도(직무 맞춤 모드)를 병렬 실행.
202+
# 첫인상·직무 적합도는 종합 점수(overall)에 미포함 — generator 가 모른 채 계산한 뒤 표시용으로 덧붙인다.
203+
(
204+
result,
205+
self_intro_item,
206+
job_fit_items,
207+
personality_item,
208+
answer_coaching,
209+
) = await asyncio.gather(
210+
_tracked(
211+
self._generate_panel(
212+
job_category=req.job_category,
213+
mode=req.mode,
214+
total_question_count=req.total_question_count,
215+
end_reason=req.end_reason,
216+
transcript=transcript,
217+
score_basis=score_basis,
218+
rag_context=rag_context,
219+
voice_analysis_summary=voice_analysis_summary,
220+
domain_question_counts=req.domain_question_counts,
221+
session_id=req.session_id,
222+
)
223+
),
224+
_tracked(self._evaluate_self_intro(req, voice_analysis_summary)),
225+
_tracked(self._evaluate_job_fit(req, transcript, rag_context)),
226+
_tracked(self._evaluate_personality(req)),
227+
_tracked(self._coach_answers(req)),
228+
)
229+
# 빈 평가위원 항목(점수·내용 모두 없음)은 표시하지 않는다 — LLM 부분 응답이 빈 패널로 새는 것 방지.
230+
extras = [self_intro_item, *job_fit_items, personality_item]
231+
result.panel_breakdown.extend(
232+
e for e in extras if e is not None and _panel_has_content(e)
233+
)
211234

212-
await self._emit_progress(
213-
session_id=req.session_id,
214-
phase="FINALIZING",
215-
message="피드백 리포트를 정리하고 있어요.",
216-
trace_id=envelope.trace_id,
217-
)
218-
payload = FeedbackCallbackPayload(
219-
session_id=req.session_id,
220-
overall_score=result.overall_score,
221-
technical_accuracy=result.technical_accuracy,
222-
logic_score=result.logic_score,
223-
communication_score=result.communication_score,
224-
strengths_summary=result.strengths_summary,
225-
weaknesses_summary=result.weaknesses_summary,
226-
improvement_keywords=result.improvement_keywords,
227-
study_plan=result.study_plan,
228-
highlights=result.highlights,
229-
panel_breakdown=result.panel_breakdown,
230-
answer_coaching=answer_coaching,
231-
report_s3_key=None,
232-
)
235+
await self._emit_progress(
236+
session_id=req.session_id,
237+
phase="FINALIZING",
238+
message="피드백 리포트를 정리하고 있어요.",
239+
trace_id=envelope.trace_id,
240+
)
241+
payload = FeedbackCallbackPayload(
242+
session_id=req.session_id,
243+
overall_score=result.overall_score,
244+
technical_accuracy=result.technical_accuracy,
245+
logic_score=result.logic_score,
246+
communication_score=result.communication_score,
247+
strengths_summary=result.strengths_summary,
248+
weaknesses_summary=result.weaknesses_summary,
249+
improvement_keywords=result.improvement_keywords,
250+
study_plan=result.study_plan,
251+
highlights=result.highlights,
252+
panel_breakdown=result.panel_breakdown,
253+
answer_coaching=answer_coaching,
254+
report_s3_key=None,
255+
)
233256

234-
await self._publisher.publish(
235-
routing_key=self._callback_routing_key,
236-
message_type="callback.feedback",
237-
payload=payload,
238-
trace_id=envelope.trace_id,
239-
correlation_id=envelope.message_id,
240-
context=envelope.context,
241-
)
242-
log.info(
243-
"feedback.generate.done",
257+
return payload
258+
259+
async def _publish_callback(
260+
self,
261+
envelope: Envelope[GenerateFeedbackRequest],
262+
payload: FeedbackCallbackPayload,
263+
) -> None:
264+
await self._publisher.publish(
265+
routing_key=self._callback_routing_key,
266+
message_type="callback.feedback",
267+
payload=payload,
268+
trace_id=envelope.trace_id,
269+
correlation_id=envelope.message_id,
270+
context=envelope.context,
271+
)
272+
273+
async def _publish_failed(
274+
self, envelope: Envelope[GenerateFeedbackRequest], exc: Exception
275+
) -> None:
276+
"""생성 중 예상 못 한 예외의 실패 신호. 콜백 없이 DLQ 로만 격리되면 Core 가 실패를
277+
모른 채 세션이 '피드백 생성 중'에 무기한 멈춘다 — 항상 FAILED 콜백을 발행하고
278+
ack 한다. 폴백 발행마저 실패하면 멱등 마킹을 되돌리고 원 예외를 다시 던져
279+
DLQ 로 보낸다(최후 안전망 — 재주입 시 duplicate skip 으로 삼켜지지 않게)."""
280+
req = envelope.payload
281+
log.exception(
282+
"feedback.generate.unexpected",
283+
message_id=envelope.message_id,
284+
session_id=req.session_id,
285+
trace_id=envelope.trace_id,
286+
)
287+
# questions/followup consumer 와 동일 분류 — TypeError(LLM 출력 스키마 불일치)는
288+
# 같은 입력으로 재시도해도 똑같이 죽는다 → retriable=false.
289+
is_schema = isinstance(exc, TypeError)
290+
payload = FeedbackCallbackPayload(
291+
session_id=req.session_id,
292+
status="FAILED",
293+
error_code="GENERATION_SCHEMA_INVALID" if is_schema else "UNEXPECTED",
294+
# str(exc) 는 LLM 응답 본문·입력 repr 까지 담길 수 있다 — 로그·와이어 크기 상한.
295+
error_message=f"{type(exc).__name__}: {exc}"[:500],
296+
retriable=not is_schema,
297+
)
298+
try:
299+
await self._publish_callback(envelope, payload)
300+
except Exception: # noqa: BLE001
301+
log.exception(
302+
"feedback.failed_callback.publish_failed",
244303
message_id=envelope.message_id,
245304
session_id=req.session_id,
246305
trace_id=envelope.trace_id,
247306
)
307+
self._idempotency.unmark(envelope.message_id)
308+
raise exc
248309

249310
async def _emit_progress(
250311
self,

ai/src/ai_server/messaging/idempotency.py

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,3 +25,12 @@ def is_seen_then_mark(self, key: str) -> bool:
2525
if len(self._store) > self._max_size:
2626
self._store.popitem(last=False)
2727
return False
28+
29+
def unmark(self, key: str) -> None:
30+
"""처리 실패로 콜백을 하나도 발행하지 못한 메시지의 마킹 해제.
31+
32+
마킹은 처리 전에 이뤄지므로, 그대로 두면 DLQ 로 간 메시지를 같은 프로세스에
33+
재주입했을 때 duplicate 로 skip 돼 콜백 없이 ack 된다 — 재처리 가능하게 되돌린다.
34+
"""
35+
with self._lock:
36+
self._store.pop(key, None)

ai/src/ai_server/model/messages/feedback.py

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
from pydantic import BaseModel, Field
44

55
from ai_server.model._config import camel_config
6+
from ai_server.model.messages.questions import GenerationStatus
67

78
InterviewMode = Literal["PERSONALITY", "TECHNICAL", "INTEGRATED", "JOB_TAILORED"]
89

@@ -120,3 +121,8 @@ class FeedbackCallbackPayload(BaseModel):
120121
# 질문별 복기(답변 메시지별 모범 답안·리라이트·코칭). 비면 복기 없음.
121122
answer_coaching: list[AnswerCoachingItem] = Field(default_factory=list)
122123
report_s3_key: str | None = None
124+
# 실패 신호 (questions/analysis 콜백과 동일 규약). status 미명시(구버전)는 OK 취급.
125+
status: GenerationStatus = "OK"
126+
error_code: str | None = None
127+
error_message: str | None = None
128+
retriable: bool | None = None

0 commit comments

Comments
 (0)