Skip to content

Commit f2d6cb7

Browse files
authored
Merge pull request #20 from 8ocket/develop
Develop
2 parents a2e3680 + e4d6645 commit f2d6cb7

9 files changed

Lines changed: 147 additions & 18 deletions

File tree

app/api/sessions.py

Lines changed: 36 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22
import uuid
33
from typing import Annotated
44

5-
from fastapi import APIRouter, Depends, Request
5+
from fastapi import APIRouter, Depends, HTTPException, Request
66
from fastapi.responses import StreamingResponse
77
from sqlalchemy.orm import Session
88

@@ -11,7 +11,7 @@
1111
from app.exceptions import MindLogError
1212
from app.limiter import limiter
1313
from app.schemas.session import MessageCreateRequest
14-
from app.services import chat_service
14+
from app.services import chat_service, session_service
1515
from app.utils.sse import format_sse_done, format_sse_error, stream_sse_events
1616

1717
logger = logging.getLogger(__name__)
@@ -82,7 +82,31 @@ async def generate():
8282
finally:
8383
yield format_sse_done()
8484

85-
return StreamingResponse(generate(), media_type="text/event-stream")
85+
return StreamingResponse(
86+
generate(),
87+
media_type="text/event-stream",
88+
headers={
89+
"Cache-Control": "no-cache",
90+
"X-Accel-Buffering": "no",
91+
"Connection": "keep-alive",
92+
},
93+
)
94+
95+
96+
@router.delete("/{session_id}", status_code=204)
97+
async def delete_session(
98+
session_id: uuid.UUID,
99+
db: DbSession,
100+
):
101+
"""
102+
세션 삭제 시 AI 소유 데이터 정리.
103+
104+
Java BE가 counseling_sessions 및 cascade 테이블 삭제 완료 후 호출한다.
105+
AI 서비스는 embeddings 테이블의 해당 session_id 레코드를 삭제한다.
106+
"""
107+
deleted = session_service.delete_session_data(session_id, db)
108+
if not deleted:
109+
raise HTTPException(status_code=404, detail="해당 세션의 AI 데이터를 찾을 수 없습니다.")
86110

87111

88112
@router.post("/{session_id}/finalize")
@@ -126,4 +150,12 @@ async def generate():
126150
finally:
127151
yield format_sse_done()
128152

129-
return StreamingResponse(generate(), media_type="text/event-stream")
153+
return StreamingResponse(
154+
generate(),
155+
media_type="text/event-stream",
156+
headers={
157+
"Cache-Control": "no-cache",
158+
"X-Accel-Buffering": "no",
159+
"Connection": "keep-alive",
160+
},
161+
)

app/models/knowledge.py

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -165,6 +165,12 @@ class Embedding(Base):
165165
nullable=False,
166166
comment="원본_ID",
167167
)
168+
session_id: Mapped[uuid.UUID | None] = mapped_column(
169+
UUID(as_uuid=True),
170+
nullable=True,
171+
index=True,
172+
comment="세션_ID (삭제용 메타데이터, FK 없음)",
173+
)
168174
user_id: Mapped[uuid.UUID] = mapped_column(
169175
UUID(as_uuid=True),
170176
ForeignKey("users.user_id", ondelete="CASCADE"),

app/models/session.py

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -52,7 +52,6 @@ class CounselingSession(Base):
5252
)
5353
summary_id: Mapped[uuid.UUID | None] = mapped_column(
5454
UUID(as_uuid=True),
55-
ForeignKey("session_summaries.summary_id"),
5655
nullable=True,
5756
comment="세션_요약_ID",
5857
)
@@ -83,7 +82,6 @@ class CounselingSession(Base):
8382
back_populates="session",
8483
uselist=False,
8584
cascade="all, delete-orphan",
86-
foreign_keys="[CounselingSession.summary_id]",
8785
)
8886
emotion_extractions: Mapped[list[EmotionExtraction]] = relationship(
8987
back_populates="session", cascade="all, delete-orphan"

app/services/chat_service.py

Lines changed: 21 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,8 @@
2323
from sqlalchemy.orm import Session
2424

2525
from app.config import get_settings
26+
from langchain_core.messages import AIMessageChunk
27+
2628
from app.graph.builder import get_graph
2729
from app.models.session import CounselingSession
2830
from app.services import redis_service
@@ -106,18 +108,25 @@ async def stream_message(
106108
output_blocked = False
107109
_t0 = time.monotonic()
108110
try:
109-
async for event in graph.astream_events(initial_state, config=config, version="v2"):
110-
if event["event"] == "on_chat_model_stream":
111-
chunk_content = event["data"]["chunk"].content
112-
if chunk_content:
113-
full_chunks.append(chunk_content)
114-
yield "chunk", {"content": chunk_content}
115-
elif event["event"] == "on_chat_model_end":
116-
output = event["data"].get("output")
117-
if output and hasattr(output, "response_metadata"):
118-
action = output.response_metadata.get("amazon-bedrock-guardrailAction")
119-
if action == "INTERVENED":
120-
output_blocked = True
111+
last_token = None
112+
async for chunk in graph.astream(
113+
initial_state, config=config, stream_mode="messages", version="v2"
114+
):
115+
if chunk["type"] != "messages":
116+
continue
117+
token, _metadata = chunk["data"]
118+
if not isinstance(token, AIMessageChunk):
119+
continue
120+
last_token = token
121+
if token.content:
122+
full_chunks.append(token.content)
123+
yield "chunk", {"content": token.content}
124+
125+
# Bedrock 가드레일 출력 차단 여부 확인 (마지막 토큰의 response_metadata)
126+
if last_token is not None and hasattr(last_token, "response_metadata"):
127+
action = last_token.response_metadata.get("amazon-bedrock-guardrailAction")
128+
if action == "INTERVENED":
129+
output_blocked = True
121130
except Exception as e:
122131
logger.error(
123132
"LangGraph 스트리밍 실패 | session_id=%s | error=%s", session_id, e, exc_info=True

app/services/embedding_service.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -133,6 +133,7 @@ async def store_session_embedding(
133133
embedding = Embedding(
134134
source_type="session_context",
135135
source_id=session_id,
136+
session_id=session_id,
136137
user_id=user_id,
137138
vector=vector,
138139
model_version=_settings.bedrock_embedding_model_id,
@@ -151,6 +152,7 @@ async def store_proposition_embeddings(
151152
*,
152153
records: list[tuple[SessionProposition, str]],
153154
user_id: uuid.UUID,
155+
session_id: uuid.UUID,
154156
session_summary_fact: str | None = None,
155157
) -> None:
156158
"""
@@ -189,6 +191,7 @@ async def store_proposition_embeddings(
189191
Embedding(
190192
source_type="proposition",
191193
source_id=prop.proposition_id,
194+
session_id=session_id,
192195
user_id=user_id,
193196
vector=vector,
194197
model_version=_settings.bedrock_embedding_model_id,

app/services/knowledge_service.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -302,6 +302,7 @@ async def extract_and_store_propositions(
302302
db,
303303
records=records,
304304
user_id=user_id,
305+
session_id=session_id,
305306
session_summary_fact=session_summary_fact,
306307
)
307308

app/services/session_service.py

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,43 @@
1+
import logging
2+
import uuid
3+
4+
from sqlalchemy.orm import Session
5+
6+
from app.models.knowledge import Embedding
7+
8+
logger = logging.getLogger(__name__)
9+
10+
11+
def delete_session_data(session_id: uuid.UUID, db: Session) -> bool:
12+
"""
13+
세션 삭제 시 AI 서비스가 자체 관리하는 데이터를 정리한다.
14+
15+
호출 시점: Java BE가 counseling_sessions 및 cascade 테이블을 삭제한 이후.
16+
17+
삭제 대상:
18+
- embeddings (source_type='session_context', source_type='proposition')
19+
→ Embedding.session_id 컬럼으로 일괄 삭제
20+
21+
처리하지 않는 항목:
22+
- session_messages, session_propositions 등: Java BE cascade로 이미 삭제됨
23+
- knowledge_relations.session_id: ondelete="SET NULL"으로 자동 null화
24+
- knowledge_entities: 사용자 지식 누적 데이터이므로 유지
25+
- Redis rag:{user_id}:* 캐시: TTL=180초 자연 만료
26+
27+
Returns:
28+
True — 임베딩 레코드가 존재해 삭제 성공
29+
False — 해당 session_id의 임베딩 레코드 없음 (404 반환용)
30+
"""
31+
deleted_count = (
32+
db.query(Embedding)
33+
.filter(Embedding.session_id == session_id)
34+
.delete(synchronize_session=False)
35+
)
36+
db.commit()
37+
38+
if deleted_count > 0:
39+
logger.info("세션 임베딩 삭제 완료 | session_id=%s | count=%d", session_id, deleted_count)
40+
else:
41+
logger.warning("삭제할 임베딩 없음 | session_id=%s", session_id)
42+
43+
return deleted_count > 0

app/utils/sse.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44
API 레이어의 중복 SSE 포맷팅 코드를 한 곳으로 통합한다.
55
"""
66

7+
import asyncio
78
import json
89
from collections.abc import AsyncIterator
910

@@ -50,3 +51,4 @@ async def stream_sse_events(
5051
yield format_sse(sse_event_name, {"content": event_data.get("content", "")})
5152
else:
5253
yield format_sse(sse_event_name, event_data)
54+
await asyncio.sleep(0) # OS TCP 버퍼 flush 기회 제공
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
"""add_session_id_to_embeddings
2+
3+
embeddings 테이블에 session_id 컬럼 추가.
4+
Java BE 선삭제 후 AI DELETE API 호출 구조에서,
5+
session_propositions가 이미 삭제된 시점에도
6+
embeddings를 session_id로 직접 삭제할 수 있도록 메타데이터 컬럼 추가.
7+
8+
Revision ID: l2m3n4o5p6q7
9+
Revises: be83bf3c61d9
10+
Create Date: 2026-04-14 00:00:00.000000
11+
"""
12+
13+
from typing import Sequence, Union
14+
15+
import sqlalchemy as sa
16+
from alembic import op
17+
from sqlalchemy.dialects.postgresql import UUID
18+
19+
revision: str = "l2m3n4o5p6q7"
20+
down_revision: Union[str, None] = "3488b28d150f"
21+
branch_labels: Union[str, Sequence[str], None] = None
22+
depends_on: Union[str, Sequence[str], None] = None
23+
24+
25+
def upgrade() -> None:
26+
op.add_column(
27+
"embeddings",
28+
sa.Column("session_id", UUID(as_uuid=True), nullable=True, comment="세션_ID (삭제용 메타데이터, FK 없음)"),
29+
)
30+
op.create_index("ix_embeddings_session_id", "embeddings", ["session_id"])
31+
32+
33+
def downgrade() -> None:
34+
op.drop_index("ix_embeddings_session_id", table_name="embeddings")
35+
op.drop_column("embeddings", "session_id")

0 commit comments

Comments
 (0)