From cc683d90c37a49c91069346e25f136c716b65392 Mon Sep 17 00:00:00 2001 From: Jaeho Date: Sat, 22 Aug 2026 19:19:15 +0900 Subject: [PATCH 1/2] =?UTF-8?q?docs(ai):=20=EB=A9=94=EC=8B=9C=EC=A7=95?= =?UTF-8?q?=C2=B7internal=20API=20=EA=B3=84=EC=95=BD=20=EB=AC=B8=EC=84=9C?= =?UTF-8?q?=EB=A5=BC=20=EC=BD=94=EB=93=9C=20=EC=A0=95=EB=B3=B8=EA=B3=BC=20?= =?UTF-8?q?=EB=8F=99=EA=B8=B0=ED=99=94?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - messaging.md §1·§4: definitions.json 에 이미 정의된 feedback/voice/tts 콜백 큐·DLQ 를 정식 표에 편입하고 '추가 예정'·'(예정)' 표기 제거 - §5.7: GeneratedQuestion 의 jobCategory·targetEvidence·expectedSignal 반영 - §5.9: FollowupCallbackPayload 에 없는 voiceAnalysis 블록 제거 - §5.9c/5.9d 신설: analyze.voice / callback.voice 스키마 문서화 (VoiceCallbackPayload 기준) - §5.11: 프론트가 소비 중인 highlights[] 추가 - §10 + api-conventions.md §2.6·§10: internal API 누락 2건 (embeddings/search, ai-logs) 추가 --- docs/api-conventions.md | 67 ++++++++++++++++++++++++++++++ docs/messaging.md | 92 +++++++++++++++++++++++++++++++++-------- 2 files changed, 142 insertions(+), 17 deletions(-) diff --git a/docs/api-conventions.md b/docs/api-conventions.md index 48b83a9b..d0eab0fc 100644 --- a/docs/api-conventions.md +++ b/docs/api-conventions.md @@ -98,6 +98,8 @@ GET /api/system/version 버전 정보 ``` GET /api/internal/users/{userId}/github-token AI → Core PUT /api/internal/documents/{documentId}/embeddings AI → Core +POST /api/internal/embeddings/search AI → Core +POST /api/internal/ai-logs AI → Core ``` 자세한 요청·응답 스키마와 인증 규약은 §10 내부 API 부록 참조. @@ -385,3 +387,68 @@ HTTP/1.1 202 Accepted **Idempotency** - 같은 `(documentId, chunkIndex)` 조합 재호출 시 row를 덮어씀 (INSERT ... ON CONFLICT UPDATE) - 부분 재시도 시 누락 청크가 발생하지 않도록 AI는 전체 청크를 한 번에 보내는 것을 권장 + +### 10.4 `POST /api/internal/embeddings/search` + +RAG 검색. AI가 질문 풀·꼬리질문·피드백 생성 시 컨텍스트 청크를 가져온다. `document_embeddings`에서 pgvector cosine topK — `queryText`가 주어지면 벡터 + full-text(BM25) RRF 하이브리드로 검색한다. + +**Request body** +```json +{ + "queryEmbedding": [0.012, -0.003, ...], + "queryText": "낙관적 락 동시성 제어", + "documentIds": [12, 13], + "topK": 5 +} +``` + +| 필드 | 타입 | 비고 | +|------|------|------| +| `queryEmbedding` | float[] | 필수. 쿼리 임베딩 벡터 | +| `queryText` | string | 선택. 주어지면 하이브리드(RRF) 검색 | +| `documentIds` | long[] | 검색 범위 제한. 비어 있으면 전체 | +| `topK` | int | 양수. 미지정 시 5 | + +**Response 200** +```json +{ + "hits": [ + { "documentId": 12, "chunkIndex": 3, "chunkText": "...", "distance": 0.18 } + ] +} +``` + +**Response 400** — `queryEmbedding` 누락 +**Response 401/403** — `INTERNAL_AUTH_FAILED` + +> RAG 보강용이라 fatal 이 아니다 — AI(`core/client.py: search_embeddings`)는 호출 실패·non-2xx 를 빈 결과로 폴백하고 검색 없이 생성을 계속한다. + +### 10.5 `POST /api/internal/ai-logs` + +AI의 LLM 호출별 토큰·지연시간을 `ai_request_logs`에 기록 (US-30 관측성). LangChain 콜백(`observability/llm_logging_callback.py`)이 호출 완료·실패 시마다 발사한다. + +**Request body** +```json +{ + "userId": 123, + "sessionId": 99, + "requestType": "generate.followup", + "modelName": "gemini-3.1-flash", + "inputTokens": 1820, + "outputTokens": 210, + "latencyMs": 1830, + "status": "SUCCEEDED", + "errorMessage": null +} +``` + +| 필드 | 타입 | 비고 | +|------|------|------| +| `requestType` | string | 필수 (NotBlank) | +| `status` | string | 필수 | +| 나머지 | — | 전부 nullable | + +**Response 202** — 기록 큐잉 (본문 없음) +**Response 401/403** — `INTERNAL_AUTH_FAILED` + +> Fire-and-forget — AI는 실패해도 raise 하지 않고 경고 로그만 남긴다(관측 실패가 본 작업을 막지 않게). diff --git a/docs/messaging.md b/docs/messaging.md index 4c39448a..9c162b50 100644 --- a/docs/messaging.md +++ b/docs/messaging.md @@ -25,9 +25,14 @@ | `ai.analyze.web` | `stackup.core-to-ai` | `analyze.web` | AI Server | | `ai.generate.questions` | `stackup.core-to-ai` | `generate.questions` | AI Server | | `ai.generate.followup` | `stackup.core-to-ai` | `generate.followup` | AI Server | +| `ai.generate.feedback` | `stackup.core-to-ai` | `generate.feedback` | AI Server | +| `ai.analyze.voice` | `stackup.core-to-ai` | `analyze.voice` | AI Server | | `ai.generate.tts` | `stackup.core-to-ai` | `generate.tts` | AI Server | | `core.callback.analysis` | `stackup.ai-to-core` | `callback.analysis` | Core Server | | `core.callback.questions` | `stackup.ai-to-core` | `callback.questions` | Core Server | +| `core.callback.feedback` | `stackup.ai-to-core` | `callback.feedback` | Core Server | +| `core.callback.voice` | `stackup.ai-to-core` | `callback.voice` | Core Server | +| `core.callback.tts` | `stackup.ai-to-core` | `callback.tts` | Core Server | | `q.realtime.session.notify` | `stackup.realtime` | `realtime.session.*` · `realtime.user.*` · `realtime.document.*` | RealTime Server | ### Dead Letter Queues (durable) @@ -43,18 +48,16 @@ | `dlq.ai.analyze.web` | `stackup.dlx` | `dlq.ai.analyze.web` | `ai.analyze.web` 처리 실패 | | `dlq.ai.generate.questions` | `stackup.dlx` | `dlq.ai.generate.questions` | `ai.generate.questions` 처리 실패 | | `dlq.ai.generate.followup` | `stackup.dlx` | `dlq.ai.generate.followup` | `ai.generate.followup` 처리 실패 | +| `dlq.ai.generate.feedback` | `stackup.dlx` | `dlq.ai.generate.feedback` | `ai.generate.feedback` 처리 실패 | +| `dlq.ai.analyze.voice` | `stackup.dlx` | `dlq.ai.analyze.voice` | `ai.analyze.voice` 처리 실패 | | `dlq.ai.generate.tts` | `stackup.dlx` | `dlq.ai.generate.tts` | `ai.generate.tts` 처리 실패 | | `dlq.core.callback.analysis` | `stackup.dlx` | `dlq.core.callback.analysis` | `core.callback.analysis` 처리 실패 | | `dlq.core.callback.questions` | `stackup.dlx` | `dlq.core.callback.questions` | `core.callback.questions` 처리 실패 | +| `dlq.core.callback.feedback` | `stackup.dlx` | `dlq.core.callback.feedback` | `core.callback.feedback` 처리 실패 | +| `dlq.core.callback.voice` | `stackup.dlx` | `dlq.core.callback.voice` | `core.callback.voice` 처리 실패 | +| `dlq.core.callback.tts` | `stackup.dlx` | `dlq.core.callback.tts` | `core.callback.tts` 처리 실패 | | `dlq.q.realtime.session.notify` | `stackup.dlx` | `dlq.q.realtime.session.notify` | `q.realtime.session.notify` 처리 실패 | -### 추가 예정 (정의 시점에 본 표 갱신) - -| 후보 Queue | 용도 | -|------------|------| -| `ai.generate.feedback` | 세션 종료 후 종합 피드백 생성 | -| `core.callback.feedback` | 피드백 콜백 | - --- ## 2. Routing Key 명명 @@ -64,7 +67,7 @@ ``` `action` ∈ `analyze | generate | callback | realtime` -`aggregate` ∈ `resume | repository | web | questions | followup | tts | analysis | feedback | session` +`aggregate` ∈ `resume | repository | cover_letter | web | questions | followup | tts | voice | analysis | feedback | session` 새 routing key 추가 시 본 패턴 유지. @@ -126,7 +129,8 @@ | 질문 풀 생성 (US-18) | `generate.questions` | `callback.questions` | `core.callback.questions` | | 꼬리질문 생성 (US-19) | `generate.followup` | `callback.questions` | `core.callback.questions` | | 질문 TTS 합성 | `generate.tts` | `callback.tts` | `core.callback.tts` | -| 피드백 생성 (US-24) | `generate.feedback` *(예정)* | `callback.feedback` *(예정)* | `core.callback.feedback` *(예정)* | +| 음성 답변 분석 (STT + 지표) | `analyze.voice` | `callback.voice` | `core.callback.voice` | +| 피드백 생성 (US-24) | `generate.feedback` | `callback.feedback` | `core.callback.feedback` | | 세션 알림 (RT2 SSE) | `realtime.session.notify` | (없음 — 단방향 push) | `q.realtime.session.notify` | > `callback.analysis` 큐는 resume/web/repo 세 use case가 공유. consumer는 `payload.targetType` 으로 분기한다. @@ -254,14 +258,26 @@ "sessionId": 99, "kind": "POOL", "questions": [ - { "category": "PROJECT_DEEP_DIVE", "question": "..." }, - { "category": "CS_FUNDAMENTAL", "question": "..." } + { + "category": "PROJECT_DEEP_DIVE", + "question": "...", + "jobCategory": "BACKEND", + "targetEvidence": "이력서: 결제 시스템에서 재고 차감 동시성 문제를 낙관적 락으로 해결", + "expectedSignal": "락 전략 선택의 트레이드오프를 DB 레벨까지 설명하는지" + }, + { "category": "CS_FUNDAMENTAL", "question": "...", "jobCategory": null, "targetEvidence": "", "expectedSignal": "" } ], "status": "OK" } } ``` +| `questions[]` 필드 | 설명 | +|------|------| +| `jobCategory` | 이 질문이 겨냥한 직군(세션 `jobCategories` 중 하나). 다직군 패널 가중에 사용. LLM 이 비우면(null) Core 가 대표 직군으로 폴백 | +| `targetEvidence` | 질문이 근거한 자료 인용 (PROJECT/TECH 는 필수, 그 외 빈 문자열 허용). 라이브 화면에 힌트로 노출 | +| `expectedSignal` | 좋은 답이 드러내야 할 것 — 내부 평가용. 꼬리질문 채점의 `parentExpectedSignal` 로 전달(§5.8). **라이브 비노출** (정답 유출 방지) | + **실패 시** (`generate()` 예외 — LLM 게이트웨이 장애, 파싱 실패 등): ```json { @@ -317,18 +333,16 @@ "answerEvaluation": { "specificity": 3.5, "logic": 4.0, - "structure": "PARTIAL_STAR" - }, - "voiceAnalysis": { - "speakingRateWpm": 142.0, - "fillerWordCounts": { "음": 5, "어": 3 }, - "silenceDurationSec": 8.2 + "structure": "PARTIAL_STAR", + "correctness": null }, "status": "OK" } } ``` +> 음성 지표는 이 콜백에 포함되지 않는다 — 별도 파이프라인 `analyze.voice` → `callback.voice` (§5.9d) 로 전달된다. + **실패 시** (생성 실패 — 스트리밍/비스트리밍 공통): ```json { @@ -391,6 +405,45 @@ > 실패 시 `status: "FAILED"` + `errorCode`(`TTS_API_ERROR`/`TTS_STORAGE_FAILED` 등), `audioKey`/`durationSec` 는 null. OpenAI TTS 는 duration 을 주지 않으므로 `durationSec` 는 null 일 수 있다. +### 5.9c `analyze.voice` +```json +{ + "messageType": "analyze.voice", + "payload": { + "sessionId": 99, + "messageId": 502, + "parentQuestionMessageId": 501, + "audioS3Key": "interview/answers/99/502.webm", + "contentType": "audio/webm", + "previousQuestionText": "...", + "mode": "TECHNICAL", + "jobCategory": "BACKEND" + }, + "context": { "userId": 123, "sessionId": 99 } +} +``` + +> Core 가 음성 답변 업로드 commit 후 발행(§5.14). `messageId` 는 `interview_messages.id` — STT 후 content 를 채울 placeholder. AI 가 S3 오디오를 받아 STT + 음성 지표(WPM/간투어/침묵) 계산 후 `callback.voice` 회신. + +### 5.9d `callback.voice` +```json +{ + "messageType": "callback.voice", + "payload": { + "sessionId": 99, + "interviewMessageId": 502, + "transcript": "네, 저는 결제 시스템에서...", + "speakingRateWpm": 142.0, + "silenceDurationSec": 8.2, + "fillerWordCounts": { "음": 5, "어": 3 }, + "pronunciationAccuracy": 0.93, + "errorCode": null + } +} +``` + +> STT 결과 + 음성 지표. 지표 필드는 전부 nullable — 계산 불가 시 null. 실패 시 `errorCode` 가 채워지고 Core 는 해당 메시지를 FAILED 로 마킹한다. 세션 종합 피드백에는 이 지표를 Core 가 집계해 `generate.feedback` 의 `voiceAnalysisSummary` 로 동봉한다(개별 콜백을 AI 가 재수집하지 않음). + ### 5.10 `generate.feedback` > `messages[]` 의 각 항목은 `category` 를 포함한다(질문 유형). AI 는 `category=SELF_INTRODUCTION` @@ -422,6 +475,8 @@ > 직무 맞춤 모드는 **`evaluator="직무 적합도"`(역량 매칭)** + **`evaluator="직무 이해도"`(직무 이해·동기)** > 항목이 추가로 포함된다 — 모두 **종합 점수(overallScore) 집계에서 제외**된 별도 정성 평가다(메인 > generator 가 모른 채 overall 계산 후 표시용으로 append). +> `highlights[]` 는 강점·개선 본문에서 발췌한 핵심 구절 — 프론트가 부분 문자열 매칭으로 리포트에 +> 하이라이트 표시한다(`HighlightedText`). 빈 리스트 허용. ```json { @@ -436,6 +491,7 @@ "weaknessesSummary": "...", "improvementKeywords": ["JPA 영속성 컨텍스트", "TCP 3-way handshake"], "studyPlan": ["..."], + "highlights": ["결론부터 말하는 답변 구조", "동시성 제어 경험"], "answerCoaching": [ { "messageId": 203, "modelAnswer": "이 질문에 강한 답변 예시…", "answerRewrite": "내 답변을 이렇게 고치면…", "coachingComment": "결론을 먼저 말하세요." } ], @@ -621,6 +677,8 @@ docker exec stackup-rabbitmq rabbitmqadmin \ |--------|------|--------|------| | `GET` | `/api/internal/users/{userId}/github-token` | AI | 사용자별 GitHub access token을 분석 시점에 짧게 위임 (envelope에 비밀 미동봉) | | `PUT` | `/api/internal/documents/{documentId}/embeddings` | AI | 청크 + 임베딩을 `document_embeddings`에 idempotent upsert | +| `POST` | `/api/internal/embeddings/search` | AI | RAG 검색 — pgvector cosine topK (queryText 동봉 시 벡터+BM25 RRF 하이브리드). 실패 시 AI 는 빈 결과로 폴백 (non-fatal) | +| `POST` | `/api/internal/ai-logs` | AI | LLM 호출별 토큰·지연시간을 `ai_request_logs` 에 기록 (fire-and-forget, 실패 무시) | 요청·응답 스키마 및 인증 규약은 [`/docs/api-conventions.md §10`](./api-conventions.md) 참조. From 60ae3caf5be7e5a179dffbaeec9ea99cf748ee62 Mon Sep 17 00:00:00 2001 From: Jaeho Date: Sat, 22 Aug 2026 19:19:20 +0900 Subject: [PATCH 2/2] =?UTF-8?q?docs(ai):=20ai=5Fserver=20=ED=8C=A8?= =?UTF-8?q?=ED=82=A4=EC=A7=80=20=EA=B0=80=EC=9D=B4=EB=93=9C=EC=9D=98=20sta?= =?UTF-8?q?le=20=EC=B0=B8=EC=A1=B0=20=EC=A0=95=EC=A0=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 존재하지 않는 chain/parsers/·rag/pgvector_client.py·storage/keys.py·analyzer/feedback_generator.py 제거 - 실제 파일명으로 교체: rag/splitter.py→chunker.py, analyzer/repo_analyzer.py→repository_analyzer.py, model/messages.py→model/messages/ 패키지, messaging/consumers/ 구조 반영 - 미기재 모듈 core/·observability/ 추가, voice '(Phase 2)'·runner '(도입 예정)' 라벨을 실구현 상태로 갱신 - §3 lifespan 예시를 MessagingRuntime 실제 패턴으로 교체 --- ai/src/ai_server/CLAUDE.md | 115 +++++++++++++++++++++---------------- 1 file changed, 66 insertions(+), 49 deletions(-) diff --git a/ai/src/ai_server/CLAUDE.md b/ai/src/ai_server/CLAUDE.md index 0e71d3b4..b23eeba0 100644 --- a/ai/src/ai_server/CLAUDE.md +++ b/ai/src/ai_server/CLAUDE.md @@ -15,10 +15,11 @@ messaging │ └──→ voice ──→ │ │ │ └──→ rag ────┴──→ storage (S3) - └──→ httpx (Core API) + └──→ core (httpx, Core internal API) config: 모두가 의존 model: 모두가 의존 (Pydantic 스키마) +observability: chain 에 붙는 LangChain 콜백 — core 경유로 호출 로그 POST ``` 원칙: @@ -45,22 +46,21 @@ settings = Settings() # singleton ``` ### `model/` -- RabbitMQ envelope 모델 -- 도메인 객체 (`AnalyzedResume`, `QuestionPool`, `FollowUpResult`) -- LLM 응답 schema (Pydantic) — `OutputParser`에 사용 +- RabbitMQ envelope 모델 (`envelope.py`) +- 메시지 페이로드 — `model/messages/` 하위에 도메인별 파일 + (`analyze.py`, `questions.py`, `followup.py`, `feedback.py`, `voice.py`, `tts.py`, `realtime.py`) +- LLM 응답 schema (Pydantic) — 체인의 구조화 출력 검증에 사용 +- 모든 페이로드 모델은 `_config.camel_config()` 로 wire 필드명을 camelCase 로 직렬화 ```python -# model/messages.py -class ResumeAnalyzeRequest(BaseModel): - resume_id: int - s3_key: str - -class ResumeAnalyzed(BaseModel): - resume_id: int - summary: str - tech_stack: list[str] - document_s3_key: str - embedding_chunk_count: int +# model/messages/questions.py (발췌) +class QuestionPoolCallbackPayload(BaseModel): + model_config = camel_config() + + session_id: int + kind: CallbackKind = "POOL" + questions: list[GeneratedQuestion] = [] + status: GenerationStatus = "OK" ``` ### `api/` @@ -70,11 +70,15 @@ class ResumeAnalyzed(BaseModel): ### `messaging/` - aio-pika consumer / publisher -- 큐별 consumer 함수 분리 +- 큐별 consumer 는 `messaging/consumers/{name}_consumer.py` 로 분리 + (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`(델타/오디오) - 모든 consumer는 envelope parsing → trace_context → 비즈니스 핸들러 호출 패턴 ```python -# messaging/resume_consumer.py +# messaging/consumers/resume_consumer.py (패턴) async def consume(message: AbstractIncomingMessage) -> None: async with message.process(requeue=False): envelope = parse_envelope(message) @@ -86,28 +90,40 @@ async def consume(message: AbstractIncomingMessage) -> None: ``` ### `analyzer/` -- use case 단위 (`resume_analyzer.py`, `repo_analyzer.py`, `feedback_generator.py`) -- 외부 입력 → 내부 모듈 조합 → 결과 publish +- 분석 use case 단위 (`resume_analyzer.py`, `repository_analyzer.py`, `web_resume_analyzer.py`) +- 소스 추출 추상화는 `analyzer/sources/` (PDF/GitHub/웹/텍스트), 임베딩 인제스트는 `_embedding_step.py` +- 외부 입력 → 내부 모듈 조합 → 결과 publish. 피드백 생성은 analyzer 가 아니라 + `messaging/consumers/feedback_consumer.py` + `chain/feedback_generation_chain.py` 에 있다 - LLM 호출 자체는 `chain/`으로 위임 ### `chain/` -- LangChain 체인 정의 +- LangChain 체인 정의 (`document_analysis_chain.py`, `question_generation_chain.py`, + `followup_generation_chain.py`, `feedback_generation_chain.py`, `pdf_vision.py`, `sentence_split.py`) - `chain/prompts/` 하위에 prompt 템플릿 (모든 프롬프트가 한 곳에) -- `chain/parsers/` 출력 파서 +- 출력 파싱은 별도 모듈 없이 각 체인 안에서 Pydantic 구조화 출력으로 검증 ### `rag/` -- 청킹 (`splitter.py`) -- 임베딩 생성 (`embedder.py`) -- 검색 어댑터 (Core API client `pgvector_client.py`) +- 청킹 (`chunker.py` — `MarkdownChunker`) +- 임베딩 생성 (`embedder.py` — provider 추상화 + Gemini/Mock 구현) +- 검색은 rag 모듈이 아니라 `core/client.py: search_embeddings` (Core `POST /api/internal/embeddings/search`) -### `voice/` (Phase 2) -- `voice/stt/` — interface + provider impls -- `voice/tts/` -- `voice/analysis/` — WPM, filler, silence +### `core/` +- Core 내부 API httpx 클라이언트 (`client.py`) — `X-Internal-API-Key` 인증 +- GitHub token 위임 · 임베딩 upsert/검색 · AI 호출 로그 기록 (엔드포인트 목록: + [`/docs/messaging.md §10`](../../../docs/messaging.md)) + +### `voice/` +- `voice/stt/` — interface + provider impls (배치 Whisper/Deepgram + 라이브 Deepgram Live) +- `voice/tts/` — provider 추상화 (Gateway/Gemini/OpenAI/Mock) +- `voice/analysis/` — WPM, filler, silence (`metrics.py`) ### `storage/` -- S3 client wrapper (`s3.py`) -- key 생성 헬퍼 (`keys.py`) — [`/docs/storage.md §2`](../../../docs/storage.md) 컨벤션 준수 +- `ObjectStorage` 추상화 (`base.py`) + `s3.py` / `local_fs.py` 구현, `factory.py` 로 토글 +- 객체 key 는 각 사용처에서 [`/docs/storage.md §2`](../../../docs/storage.md) 컨벤션대로 조립 (전용 헬퍼 모듈 없음) + +### `observability/` +- `llm_logging_callback.py` — LangChain `AsyncCallbackHandler`. 토큰/latency 측정 후 + `core/client.py: record_ai_log` 로 Core `POST /api/internal/ai-logs` (fire-and-forget) --- @@ -115,23 +131,24 @@ async def consume(message: AbstractIncomingMessage) -> None: ### REST (FastAPI) - `api/health.py` — 헬스체크 -- `api/internal/*` — Core가 호출할 수 있는 동기 endpoint (필요 시) +- `api/voice_stream.py` — `/internal/voice/stream` WS (RealTime 이 프록시한 실시간 음성 답변, RT3) ### MQ Consumer -- `messaging/runner.py` (도입 예정) — 모든 consumer를 시작하는 entry -- `main.py` lifespan에서 자동 시작 (또는 별도 프로세스로 분리 검토) +- `messaging/runner.py` — `MessagingRuntime` 이 의존성(체인·스토리지·Core 클라이언트·notifier)을 + 조립하고 모든 consumer 를 시작/종료하는 단일 entry +- `main.py` lifespan 에서 `runtime.start()` / `runtime.stop()` 호출 ```python -# main.py 의 lifespan +# main.py 의 lifespan (실제 패턴) @asynccontextmanager async def lifespan(app: FastAPI): - connection = await connect_robust(settings.rabbitmq_url) - channel = await connection.channel() - await start_resume_consumer(channel) - await start_repo_consumer(channel) - await start_session_consumer(channel) - yield - await connection.close() + runtime = MessagingRuntime(settings) + app.state.messaging = runtime + try: + await runtime.start() + yield + finally: + await runtime.stop() ``` --- @@ -165,25 +182,25 @@ class AnalysisError(Exception): 이력서 분석 (US-09)을 예로 들면: -1. `model/messages.py`에 `ResumeAnalyzeRequest`, `ResumeAnalyzed`, `ResumeFailed` 정의 -2. `messaging/resume_consumer.py` 구현 (envelope parse → handler 호출) +1. `model/messages/{name}.py`에 `ResumeAnalyzeRequest`, `ResumeAnalyzed`, `ResumeFailed` 정의 +2. `messaging/consumers/resume_consumer.py` 구현 (envelope parse → handler 호출) 3. `analyzer/resume_analyzer.py` 구현 ```python async def handle(req: ResumeAnalyzeRequest) -> None: - pdf_bytes = await s3.get(req.s3_key) + pdf_bytes = await storage.get(req.s3_key) text = extract_text(pdf_bytes) result = await resume_chain.ainvoke({"text": text}) md_key = f"analyzed/resume/{req.resume_id}/summary.md" - await s3.put(md_key, result.markdown) - chunks = split(result.markdown) + await storage.put(md_key, result.markdown) + chunks = chunker.split(result.markdown) embeddings = await embedder.embed(chunks) - await pgvector_client.upsert(req.resume_id, chunks, embeddings) + await core_client.upsert_embeddings(document_id=req.analyzed_document_id, ...) await publisher.publish_callback(ResumeAnalyzed(...)) ``` -4. `chain/resume_analyzer_chain.py` (prompt + LLM + parser) +4. `chain/{name}_chain.py` (prompt + LLM + Pydantic 구조화 출력) 5. 단위 테스트 (mock LLM) 6. 통합 테스트 (Testcontainer RabbitMQ + MinIO) -7. main.py lifespan에 consumer 등록 +7. `messaging/runner.py` 의 `MessagingRuntime` 에 consumer 등록 ---