Spaces:
Sleeping
Sleeping
File size: 8,117 Bytes
fc73f93 |
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 |
# ============================================================================
# backend/app/services/streaming_service.py - NEW FILE
# ============================================================================
"""
Streaming Service - Server-Sent Events (SSE)
Handles real-time streaming of AI responses.
Integrates with chat_service.py RAG pipeline.
"""
import asyncio
import json
from typing import AsyncGenerator, Dict, Any, List, Optional
from datetime import datetime
from app.config import settings
from app.ml.policy_network import predict_policy_action
from app.ml.retriever import retrieve_documents, format_context
from app.core.llm_manager import llm_manager
# ============================================================================
# STREAMING SERVICE
# ============================================================================
class StreamingService:
"""
Handles SSE streaming for real-time chat responses.
Events sent:
- status: Progress updates (retrieval, generation stages)
- content: Response chunks (word by word)
- metadata: Final stats (policy action, docs retrieved, etc.)
- done: Stream completion signal
- error: Error occurred
"""
def __init__(self):
print("🌊 StreamingService initialized")
async def stream_chat_response(
self,
query: str,
conversation_history: List[Dict[str, str]] = None,
user_id: Optional[str] = None
) -> AsyncGenerator[str, None]:
"""
Stream chat response with progress updates.
Yields SSE-formatted events:
- event: status, content, metadata, done, error
- data: JSON payload
Args:
query: User query
conversation_history: Previous messages
user_id: User ID
Yields:
str: SSE formatted events
"""
import time
start_time = time.time()
if conversation_history is None:
conversation_history = []
try:
# ================================================================
# STAGE 1: Policy Decision
# ================================================================
yield self._format_sse_event(
event="status",
data={"stage": "policy", "message": "Analyzing query..."}
)
await asyncio.sleep(0.1) # Small delay for UX
policy_result = predict_policy_action(
query=query,
history=conversation_history,
return_probs=True
)
# ================================================================
# STAGE 2: Retrieval (if needed)
# ================================================================
retrieved_docs = []
context = ""
retrieval_time = 0
if policy_result['should_retrieve']:
yield self._format_sse_event(
event="status",
data={"stage": "retrieval", "message": "Searching knowledge base..."}
)
retrieval_start = time.time()
try:
retrieved_docs = retrieve_documents(
query=query,
top_k=settings.TOP_K,
min_similarity=settings.SIMILARITY_THRESHOLD
)
retrieval_time = (time.time() - retrieval_start) * 1000
if retrieved_docs:
context = format_context(
retrieved_docs,
max_context_length=settings.MAX_CONTEXT_LENGTH
)
yield self._format_sse_event(
event="status",
data={
"stage": "retrieval",
"message": f"Found {len(retrieved_docs)} relevant documents"
}
)
except Exception as e:
print(f"⚠️ Retrieval error during streaming: {e}")
# Continue without retrieval
# ================================================================
# STAGE 3: Stream Generation
# ================================================================
yield self._format_sse_event(
event="status",
data={"stage": "generation", "message": "Generating response..."}
)
generation_start = time.time()
full_response = ""
# Stream from LLM
async for chunk in llm_manager.stream_chat_response(
query=query,
context=context,
history=conversation_history
):
full_response += chunk
yield self._format_sse_event(
event="content",
data={"text": chunk}
)
generation_time = (time.time() - generation_start) * 1000
total_time = (time.time() - start_time) * 1000
# ================================================================
# STAGE 4: Send Metadata
# ================================================================
metadata = {
"policy_action": policy_result['action'],
"policy_confidence": policy_result['confidence'],
"documents_retrieved": len(retrieved_docs),
"top_doc_score": retrieved_docs[0]['score'] if retrieved_docs else None,
"retrieval_time_ms": round(retrieval_time, 2),
"generation_time_ms": round(generation_time, 2),
"total_time_ms": round(total_time, 2),
"timestamp": datetime.now().isoformat()
}
# Add retrieved docs metadata
if retrieved_docs:
metadata['retrieved_docs_metadata'] = [
{
'faq_id': doc['faq_id'],
'score': doc['score'],
'category': doc['category'],
'rank': doc['rank']
}
for doc in retrieved_docs
]
yield self._format_sse_event(
event="metadata",
data=metadata
)
# ================================================================
# STAGE 5: Done
# ================================================================
yield self._format_sse_event(
event="done",
data={"message": "Stream completed"}
)
except Exception as e:
print(f"❌ Streaming error: {e}")
import traceback
traceback.print_exc()
yield self._format_sse_event(
event="error",
data={"error": str(e), "message": "An error occurred during streaming"}
)
def _format_sse_event(self, event: str, data: Dict[str, Any]) -> str:
"""
Format data as SSE event.
SSE format:
event: <event_name>
data: <json_data>
(blank line to separate events)
"""
json_data = json.dumps(data, ensure_ascii=False)
return f"event: {event}\ndata: {json_data}\n\n"
# ============================================================================
# GLOBAL INSTANCE
# ============================================================================
streaming_service = StreamingService() |