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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 7 additions & 1 deletion ai/CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -397,7 +397,13 @@ docker run --env-file .env -p 8000:8000 stackup-ai
payload 반환)와 `_failed_payload(req, exc)` 팩토리만 구현한다. `errorCode` 분류는 consumer 별
계약 유지 — questions/followup: `GENERATION_SCHEMA_INVALID`(TypeError, retriable=false) |
`GENERATION_FAILED`, feedback: 동일 스키마 코드 | `UNEXPECTED`. `errorMessage` 는 공용
`format_error_message`(`ExcType: msg`, 500자 상한). Core 쪽 처리는
`format_error_message`(`ExcType: msg`, 500자 상한). **F5/F6 확장**: 분석 4개
(resume/web/repository/cover_letter)도 같은 가드로 전환(도메인 에러의 코드·retriable 분류는
각 `_failed_payload` 팩토리가 보존, `action="analyze"` 로 기존 로그 키 유지), voice/tts 는
직접-발행 모델을 유지하되 마킹 이후 전 구간을 `unmark_on_error` 로 감쌌다(재주입은
합성 재과금·세그먼트 재전송을 동반 — 수동 복구 전용). feedback 은 요청의
`attemptId` 를 성공/FAILED 콜백에 에코 — Core 가 대체된 이전 시도의 지연 FAILED 를 드롭하는
근거(messaging.md §5.10/§5.11). Core 쪽 처리는
[`backend/CLAUDE.md`](../backend/CLAUDE.md) 참고. (질문 풀 RAG `questions_rag_timeout_sec`
1.5s 하드 타임아웃과 `llm_pro_timeout_sec` 30s/`llm_flash_timeout_sec` 10s 요청 타임아웃은
이전 리팩터에서 도입되어 유지된다.)
Expand Down
10 changes: 6 additions & 4 deletions ai/src/ai_server/CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -76,10 +76,12 @@ class QuestionPoolCallbackPayload(BaseModel):
콜백 발행은 `publisher.py`, 멱등은 `idempotency.py`(`LruIdempotencyStore`),
RealTime 직접 발행은 `progress.py`(분석 진행, user 채널)·`session_notify.py`(델타/오디오/질문 풀·피드백 생성 진행, 세션 채널)
- 모든 consumer는 envelope parsing → trace_context → 비즈니스 핸들러 호출 패턴
- 생성 계열 3개(questions/followup/feedback)는 공용 가드 `consumers/failure_signal.py:
consume_with_failure_signal` 경유 — consumer 는 `_process(envelope)`(성공 payload 반환)와
`_failed_payload(req, exc)` 팩토리만 구현하고, 파싱→멱등→전 구간 가드→FAILED 콜백/성공 발행→
unmark 는 가드가 책임진다 ([`/docs/messaging.md §6`](../../../docs/messaging.md) AI Server 절)
- 생성 3개(questions/followup/feedback) + 분석 4개(resume/web/repository/cover_letter)는 공용 가드
`consumers/failure_signal.py: consume_with_failure_signal` 경유 — consumer 는
`_process(envelope)`(성공 payload 반환)와 `_failed_payload(req, exc)` 팩토리만 구현하고,
파싱→멱등→전 구간 가드→FAILED 콜백/성공 발행→unmark 는 가드가 책임진다. voice/tts 는
직접-발행 모델 유지 + 마킹 이후 전 구간을 `unmark_on_error` 로 래핑
([`/docs/messaging.md §6`](../../../docs/messaging.md) AI Server 절)

```python
# messaging/consumers/resume_consumer.py (패턴)
Expand Down
144 changes: 48 additions & 96 deletions ai/src/ai_server/messaging/consumers/cover_letter_consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,11 @@
from aio_pika.abc import AbstractIncomingMessage

from ai_server.analyzer.resume_analyzer import ResumeAnalyzeError, ResumeAnalyzer
from ai_server.messaging.consumers.failure_signal import (
analysis_done_fields,
analysis_failed_payload,
consume_with_failure_signal,
)
from ai_server.messaging.idempotency import LruIdempotencyStore
from ai_server.messaging.progress import AnalysisProgressNotifier
from ai_server.messaging.publisher import CallbackPublisher
Expand Down Expand Up @@ -35,112 +40,49 @@ def __init__(
self._progress = progress_notifier

async def handle(self, message: AbstractIncomingMessage) -> None:
async with message.process(requeue=False):
try:
envelope = Envelope[CoverLetterAnalyzeRequest].model_validate_json(
message.body
)
except Exception as exc: # parse error → DLQ-ready (auto NACK on raise)
log.error(
"cover_letter.parse.failed",
error=str(exc),
delivery_tag=message.delivery_tag,
)
raise

if self._idempotency.is_seen_then_mark(envelope.message_id):
log.info(
"cover_letter.idempotent.skip",
message_id=envelope.message_id,
trace_id=envelope.trace_id,
)
return

req = envelope.payload
log.info(
"cover_letter.analyze.start",
message_id=envelope.message_id,
cover_letter_id=req.cover_letter_id,
trace_id=envelope.trace_id,
)

payload = await self._run_and_build_payload(
req, envelope.trace_id, user_id=envelope.context.user_id
)

await self._publisher.publish(
routing_key=self._callback_routing_key,
message_type="callback.analysis",
payload=payload,
trace_id=envelope.trace_id,
correlation_id=envelope.message_id,
context=envelope.context,
)
log.info(
"cover_letter.analyze.done",
message_id=envelope.message_id,
cover_letter_id=req.cover_letter_id,
status=payload.status,
trace_id=envelope.trace_id,
)
await consume_with_failure_signal(
message,
domain="cover_letter",
action="analyze",
envelope_type=Envelope[CoverLetterAnalyzeRequest],
idempotency=self._idempotency,
publisher=self._publisher,
routing_key=self._callback_routing_key,
message_type="callback.analysis",
process=self._process,
failed_payload=self._failed_payload,
done_fields=analysis_done_fields,
expected_errors=(ResumeAnalyzeError,),
)

async def _run_and_build_payload(
self,
req: CoverLetterAnalyzeRequest,
trace_id: str,
*,
user_id: int | None,
async def _process(
self, envelope: Envelope[CoverLetterAnalyzeRequest]
) -> AnalysisCallbackPayload:
req = envelope.payload
log.info(
"cover_letter.analyze.start",
message_id=envelope.message_id,
cover_letter_id=req.cover_letter_id,
trace_id=envelope.trace_id,
)
progress = (
self._progress.emitter_for(
user_id=user_id,
user_id=envelope.context.user_id,
target_type="COVER_LETTER",
target_id=req.cover_letter_id,
trace_id=trace_id,
trace_id=envelope.trace_id,
)
if self._progress is not None
else None
)
try:
# ResumeAnalyzer 의 resume_id/file_path 는 각각 식별자/추출 locator 로 일반화돼 있어
# 자소서는 cover_letter_id 와 inline content 를 그대로 넘긴다.
result = await self._analyzer.analyze(
resume_id=req.cover_letter_id,
file_path=req.content,
analyzed_document_id=req.analyzed_document_id,
progress=progress,
)
except ResumeAnalyzeError as err:
log.warning(
"cover_letter.analyze.domain_failed",
cover_letter_id=req.cover_letter_id,
code=err.code,
retriable=err.retriable,
trace_id=trace_id,
)
return AnalysisCallbackPayload(
target_type="COVER_LETTER",
target_id=req.cover_letter_id,
status="FAILED",
error_code=err.code,
error_message=err.message,
retriable=err.retriable,
)
except Exception as exc:
log.exception(
"cover_letter.analyze.unexpected_failed",
cover_letter_id=req.cover_letter_id,
trace_id=trace_id,
)
return AnalysisCallbackPayload(
target_type="COVER_LETTER",
target_id=req.cover_letter_id,
status="FAILED",
error_code="UNEXPECTED",
error_message=str(exc),
retriable=True,
)

# ResumeAnalyzer 의 resume_id/file_path 는 각각 식별자/추출 locator 로 일반화돼 있어
# 자소서는 cover_letter_id 와 inline content 를 그대로 넘긴다.
result = await self._analyzer.analyze(
resume_id=req.cover_letter_id,
file_path=req.content,
analyzed_document_id=req.analyzed_document_id,
progress=progress,
)
return AnalysisCallbackPayload(
target_type="COVER_LETTER",
target_id=req.cover_letter_id,
Expand All @@ -150,3 +92,13 @@ async def _run_and_build_payload(
document_path=result.document_path,
embedding_chunk_count=result.embedding_chunk_count,
)

def _failed_payload(
self, req: CoverLetterAnalyzeRequest, exc: Exception
) -> AnalysisCallbackPayload:
return analysis_failed_payload(
target_type="COVER_LETTER",
target_id=req.cover_letter_id,
exc=exc,
domain_error=ResumeAnalyzeError,
)
88 changes: 74 additions & 14 deletions ai/src/ai_server/messaging/consumers/failure_signal.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
from __future__ import annotations

from collections.abc import Awaitable, Callable
import contextlib
from collections.abc import AsyncIterator, Awaitable, Callable
from typing import Any, TypeVar

import structlog
Expand All @@ -10,6 +11,7 @@
from ai_server.messaging.idempotency import LruIdempotencyStore
from ai_server.messaging.publisher import CallbackPublisher
from ai_server.model.envelope import Envelope
from ai_server.model.messages.analyze import AnalysisCallbackPayload

log = structlog.get_logger(__name__)

Expand All @@ -30,6 +32,54 @@ def classify_failure(exc: Exception, *, unexpected_code: str) -> tuple[str, bool
return unexpected_code, True


@contextlib.asynccontextmanager
async def unmark_on_error(
idempotency: LruIdempotencyStore, message_id: str
) -> AsyncIterator[None]:
"""멱등 마킹 이후 구간의 어떤 예외든 unmark 후 전파 — 발행 실패뿐 아니라 그 앞의
합성·계산 단계에서 죽어도, 콜백 0건으로 DLQ 에 간 메시지의 재주입이 duplicate skip 으로
삼켜지지 않게 한다. 전면 가드 전환이 부담스러운 직접-발행 모델 컨슈머(voice/tts)용.
주의: 재주입은 pre-publish 부수효과(합성·STT 재과금, 라이브 세그먼트 재전송)를 재실행한다 —
운영자 수동 복구 수단으로만 쓴다 (docs/messaging.md §6)."""
try:
yield
except Exception:
idempotency.unmark(message_id)
raise


def analysis_failed_payload(
*,
target_type: str,
target_id: int,
exc: Exception,
domain_error: type[Exception],
) -> AnalysisCallbackPayload:
"""분석 4종(resume/web/repository/cover_letter) 공용 FAILED payload —
도메인 에러는 코드·retriable 분류를 보존하고, 그 외는 UNEXPECTED/retriable=true."""
if isinstance(exc, domain_error):
return AnalysisCallbackPayload(
target_type=target_type,
target_id=target_id,
status="FAILED",
error_code=exc.code, # type: ignore[attr-defined]
error_message=exc.message, # type: ignore[attr-defined]
retriable=exc.retriable, # type: ignore[attr-defined]
)
return AnalysisCallbackPayload(
target_type=target_type,
target_id=target_id,
status="FAILED",
error_code="UNEXPECTED",
error_message=format_error_message(exc),
retriable=True,
)


def analysis_done_fields(payload: Any) -> dict[str, Any]:
return {"target_id": payload.target_id, "status": payload.status}


async def consume_with_failure_signal(
message: AbstractIncomingMessage,
*,
Expand All @@ -42,6 +92,8 @@ async def consume_with_failure_signal(
process: Callable[[Envelope[ReqT]], Awaitable[BaseModel]],
failed_payload: Callable[[ReqT, Exception], BaseModel],
done_fields: Callable[[Any], dict[str, Any]] | None = None,
action: str = "generate",
expected_errors: tuple[type[Exception], ...] = (),
) -> None:
"""생성 컨슈머(질문 풀·꼬리질문·피드백) 공용 실패 신호 가드.

Expand Down Expand Up @@ -84,17 +136,29 @@ async def _publish(payload: BaseModel) -> None:
context=envelope.context,
)

log_ids: dict[str, Any] = {
"message_id": envelope.message_id,
"trace_id": envelope.trace_id,
}
session_id = getattr(envelope.payload, "session_id", None)
if session_id is not None:
log_ids["session_id"] = session_id

try:
payload = await process(envelope)
except Exception as exc: # noqa: BLE001
log.exception(
f"{domain}.generate.failed",
message_id=envelope.message_id,
session_id=session_id,
trace_id=envelope.trace_id,
)
if expected_errors and isinstance(exc, expected_errors):
# 예상된 도메인 실패(빈 PDF·404 URL 등 일상 입력)는 traceback 없는 warning —
# ERROR 레벨 스택트레이스로 일상 실패가 알람을 울리지 않게 한다.
log.warning(
f"{domain}.{action}.failed",
error_code=getattr(exc, "code", None),
retriable=getattr(exc, "retriable", None),
error=str(exc),
**log_ids,
)
else:
log.exception(f"{domain}.{action}.failed", **log_ids)
try:
# 팩토리 자체가 죽어도(검증 오류 등) 같은 안전망을 태운다 —
# try 밖이면 unmark 없이 DLQ 로 가서 재주입이 duplicate skip 으로 삼켜진다.
Expand All @@ -103,9 +167,7 @@ async def _publish(payload: BaseModel) -> None:
except Exception: # noqa: BLE001
log.exception(
f"{domain}.failed_callback.publish_failed",
message_id=envelope.message_id,
session_id=session_id,
trace_id=envelope.trace_id,
**log_ids,
)
idempotency.unmark(envelope.message_id)
raise exc
Expand All @@ -117,9 +179,7 @@ async def _publish(payload: BaseModel) -> None:
idempotency.unmark(envelope.message_id)
raise
log.info(
f"{domain}.generate.done",
message_id=envelope.message_id,
session_id=session_id,
trace_id=envelope.trace_id,
f"{domain}.{action}.done",
**log_ids,
**(done_fields(payload) if done_fields is not None else {}),
)
2 changes: 2 additions & 0 deletions ai/src/ai_server/messaging/consumers/feedback_consumer.py
Original file line number Diff line number Diff line change
Expand Up @@ -228,6 +228,7 @@ async def _tracked(coro: Awaitable[T]) -> T:
panel_breakdown=result.panel_breakdown,
answer_coaching=answer_coaching,
report_s3_key=None,
attempt_id=req.attempt_id,
)

return payload
Expand All @@ -245,6 +246,7 @@ def _failed_payload(
error_code=error_code,
error_message=format_error_message(exc),
retriable=retriable,
attempt_id=req.attempt_id,
)

async def _emit_progress(
Expand Down
Loading
Loading