Repository navigation
Expand file tree
/
Copy pathmain.py
More file actions
589 lines (509 loc) Β· 19.2 KB
/
Copy pathmain.py
File metadata and controls
589 lines (509 loc) Β· 19.2 KB
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
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
"""
Test the RAG - Advanced RAG Knowledge Processing System
A professional-grade Retrieval-Augmented Generation (RAG) testing platform built with FastAPI,
featuring intelligent knowledge ingestion, vector storage, and AI-powered query processing.
This module provides the main FastAPI application with endpoints for:
- Knowledge ingestion and management
- RAG-powered query processing
- Direct LLM interactions
- System validation and monitoring
- Comprehensive RAG testing and benchmarking
Author: Professional Development Team
Version: 1.0.0
"""
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
import time
import logging
from llmhandler import llm_handler
from RAGPipeline import RAGPipeline
from KnowledgeIngestion import KnowledgeIngestion
from ConversationManager import ConversationManager
from Configurator import get_config
# Configure logging
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
handlers=[
logging.StreamHandler(),
logging.FileHandler('rag_system.log')
]
)
logger = logging.getLogger(__name__)
# Initialize FastAPI application with comprehensive metadata
app = FastAPI(
title="Test the RAG - RAG Knowledge Processing System",
description="Advanced Retrieval-Augmented Generation testing platform for intelligent knowledge processing and RAG optimization",
version="1.0.0",
docs_url="/docs",
redoc_url="/redoc"
)
# Initialize core system components
rag_pipeline = RAGPipeline()
knowledge_ingestion = KnowledgeIngestion()
conversation_manager = ConversationManager(get_config())
class InputRequest(BaseModel):
input: str
class OutputResponse(BaseModel):
output: str
class KnowledgeResponse(BaseModel):
success: bool
message: str = ""
files_processed: int = 0
chunks_created: int = 0
points_stored: int = 0
embedding_model: str = "unknown"
chunking_strategy: str = "unknown"
error: str = None
class RAGRequest(BaseModel):
query: str
class RAGResponse(BaseModel):
answer: str
sources: list
used_model: dict
latency_ms: float
total_chunks_considered: int
reranked_chunks: int = 0
context_tokens: int = 0
timing: dict = {}
error: str = None
@app.post("/process", response_model=OutputResponse, tags=["LLM Processing"])
async def process_input(request: InputRequest):
"""
Direct LLM processing endpoint.
Processes input text directly through the configured LLM without knowledge retrieval.
This endpoint bypasses the RAG pipeline for direct AI interactions.
Args:
request (InputRequest): Input text to process
Returns:
OutputResponse: LLM-generated response
Raises:
HTTPException: If processing fails
Example:
```json
{
"input": "Explain quantum computing in simple terms"
}
```
"""
logger.info(f"π [API] /process endpoint called with input length: {len(request.input)} chars")
logger.info(f"π [API] Input preview: {request.input[:100]}...")
try:
start_time = time.time()
logger.info(f"π€ [LLM] Starting direct LLM processing...")
llm_response = await llm_handler.process_request(request.input)
processing_time = (time.time() - start_time) * 1000
logger.info(f"β
[LLM] Direct processing completed in {processing_time:.2f}ms")
logger.info(f"π [LLM] Response length: {len(llm_response)} chars")
logger.info(f"π― [LLM] Model used: {llm_handler.deployment_name}")
return OutputResponse(output=llm_response)
except Exception as e:
raise HTTPException(status_code=500, detail=f"LLM processing failed: {str(e)}")
@app.get("/FeedKnowledge", response_model=KnowledgeResponse, tags=["Knowledge Management"])
async def feed_knowledge():
"""
Knowledge ingestion API endpoint.
Processes all knowledge files from the Knowledge folder through a complete pipeline:
1. Reads all .md files from the Knowledge directory
2. Chunks content using the configured chunking strategy
3. Creates embeddings using the configured embedding model
4. Stores vectors in Qdrant vector database
5. Returns comprehensive ingestion statistics
This endpoint is the primary way to populate the knowledge base for RAG queries.
Returns:
KnowledgeResponse: Detailed ingestion results including:
- Success status and message
- File processing statistics
- Chunk and vector creation counts
- Model and strategy information
- Error details if applicable
Raises:
HTTPException: If ingestion fails catastrophically
Example Response:
```json
{
"success": true,
"message": "Knowledge ingestion completed successfully",
"files_processed": 5,
"chunks_created": 23,
"points_stored": 23,
"embedding_model": "EMB_3_LARGE",
"chunking_strategy": "SEMANTIC"
}
```
"""
logger.info(f"π [API] /FeedKnowledge endpoint called - Starting knowledge ingestion")
try:
logger.info(f"π [INGESTION] Initiating complete knowledge ingestion pipeline...")
result = await knowledge_ingestion.ingest_all_knowledge()
logger.info(f"β
[INGESTION] Knowledge ingestion completed successfully")
logger.info(f"π [INGESTION] Files processed: {result.get('files_processed', 0)}")
logger.info(f"π [INGESTION] Chunks created: {result.get('chunks_created', 0)}")
logger.info(f"π [INGESTION] Points stored: {result.get('points_stored', 0)}")
logger.info(f"π― [INGESTION] Model: {result.get('embedding_model', 'unknown')}")
logger.info(f"π― [INGESTION] Strategy: {result.get('chunking_strategy', 'unknown')}")
return KnowledgeResponse(**result)
except Exception as e:
logger.error(f"β [API] /FeedKnowledge endpoint failed: {str(e)}")
return KnowledgeResponse(
success=False,
message="Knowledge ingestion failed",
files_processed=0,
chunks_created=0,
points_stored=0,
embedding_model="unknown",
chunking_strategy="unknown",
error=str(e)
)
@app.post("/rag", response_model=RAGResponse, tags=["RAG Query System"])
async def rag_query(request: RAGRequest):
"""
RAG (Retrieval-Augmented Generation) API endpoint.
Processes natural language queries using the complete RAG pipeline:
1. Embeds the query using the configured embedding model
2. Searches for similar chunks in the vector database
3. Re-ranks results using the configured re-ranking method
4. Assembles context from selected chunks
5. Generates answer using the LLM with retrieved context
This is the primary endpoint for intelligent question-answering over your knowledge base.
Args:
request (RAGRequest): Natural language query
Returns:
RAGResponse: Complete RAG response including:
- AI-generated answer
- Source documents with citations
- Model information used
- Processing latency
- Number of chunks considered
- Error details if applicable
Raises:
HTTPException: If RAG processing fails catastrophically
Example Request:
```json
{
"query": "What are the delivery guidelines for blocks?"
}
```
Example Response:
```json
{
"answer": "The delivery guidelines specify...",
"sources": [
{
"source": "delivery_guide.md",
"chunk_index": 2,
"snippet": "Delivery guidelines state that...",
"score": 0.95
}
],
"used_model": {
"embedding": "EMB_3_LARGE",
"llm": "gpt-4o"
},
"latency_ms": 2341.5,
"total_chunks_considered": 10
}
```
"""
logger.info(f"π [API] /rag endpoint called with query length: {len(request.query)} chars")
logger.info(f"π [API] Query preview: {request.query[:100]}...")
try:
logger.info(f"π [RAG] Starting RAG pipeline processing...")
result = await rag_pipeline.process_query(request.query)
logger.info(f"β
[RAG] RAG pipeline completed successfully")
logger.info(f"π [RAG] Answer length: {len(result.get('answer', ''))} chars")
logger.info(f"π [RAG] Sources found: {len(result.get('sources', []))}")
logger.info(f"β±οΈ [RAG] Latency: {result.get('latency_ms', 0):.2f}ms")
logger.info(f"π [RAG] Chunks considered: {result.get('total_chunks_considered', 0)}")
if 'conversation' in result and result['conversation'].get('enabled'):
logger.info(f"π¬ [RAG] Conversation enabled - Serial: {result['conversation'].get('serial_number')}")
logger.info(f"π¬ [RAG] Total conversations: {result['conversation'].get('total_conversations')}")
return RAGResponse(**result)
except Exception as e:
logger.error(f"β [API] /rag endpoint failed: {str(e)}")
return RAGResponse(
answer=f"Error processing RAG query: {str(e)}",
sources=[],
used_model={"embedding": "unknown", "llm": "unknown"},
latency_ms=0.0,
total_chunks_considered=0,
error=str(e)
)
@app.get("/rag/validate", tags=["System Monitoring"])
async def validate_rag_pipeline():
"""
Validate the RAG pipeline components.
Performs comprehensive validation of all RAG pipeline components:
- Embedder configuration and connectivity
- Retriever connection to Qdrant vector database
- LLM handler configuration
- Collection status and data availability
This endpoint is essential for system health monitoring and troubleshooting.
Returns:
Dict: Validation results including:
- Overall validation status
- Component-specific validation results
- Error messages and warnings
- Collection information
Example Response:
```json
{
"valid": true,
"components": {
"embedder": {
"valid": true,
"model": "EMB_3_LARGE",
"vector_size": 3072
},
"retriever": {
"valid": true,
"collection_info": {
"points_count": 150,
"status": "green"
}
},
"llm_handler": {
"valid": true,
"model": "gpt-4o"
}
},
"warnings": [],
"errors": []
}
```
"""
try:
validation = await rag_pipeline.validate_pipeline()
return validation
except Exception as e:
return {
"valid": False,
"error": str(e),
"components": {},
"errors": [str(e)],
"warnings": []
}
@app.get("/rag/info", tags=["System Monitoring"])
async def get_rag_info():
"""
Get RAG pipeline configuration information.
Returns comprehensive information about the current RAG pipeline configuration:
- Embedding model and settings
- Chunking strategy and parameters
- RAG-specific configuration (top-k, re-ranking, etc.)
- LLM model and parameters
- Vector database settings
Useful for understanding current system configuration and debugging.
Returns:
Dict: Configuration information including:
- Embedding configuration
- Chunking configuration
- RAG parameters
- LLM settings
- Vector database information
Example Response:
```json
{
"embedding": {
"model": "EMB_3_LARGE",
"vector_size": 3072,
"deployment_name": "text-embedding-3-large"
},
"chunking": {
"strategy": "SEMANTIC",
"max_chars": 1000,
"overlap_chars": 100
},
"rag": {
"top_k": 10,
"final_chunks": 5,
"rerank_method": "MMR",
"context_max_tokens": 1500
},
"llm": {
"model": "gpt-4o",
"max_tokens": 1000,
"temperature": 0.7
}
}
```
"""
try:
info = rag_pipeline.get_pipeline_info()
return info
except Exception as e:
return {"error": str(e)}
@app.get("/knowledge/status", tags=["Knowledge Management"])
async def get_knowledge_status():
"""
Get the current status of knowledge ingestion.
Provides comprehensive status information about the knowledge base:
- Knowledge files available and their metadata
- Qdrant collection status and point counts
- Current configuration settings
- System readiness for RAG queries
Essential for monitoring knowledge base health and understanding data availability.
Returns:
Dict: Status information including:
- Knowledge files information
- Vector database collection status
- Current configuration
- System readiness indicators
Example Response:
```json
{
"knowledge_files": {
"total_files": 5,
"file_names": ["guide1.md", "guide2.md"],
"total_content_length": 15420
},
"qdrant_collection": {
"name": "knowledge_embeddings",
"exists": true,
"points_count": 150
},
"configuration": {
"embedding_model": "EMB_3_LARGE",
"chunking_strategy": "SEMANTIC",
"rag_top_k": 10
}
}
```
"""
try:
status = await knowledge_ingestion.get_ingestion_status()
return status
except Exception as e:
return {"error": str(e)}
@app.delete("/knowledge/clear", tags=["Knowledge Management"])
async def clear_knowledge_base():
"""
Clear all data from the knowledge base.
Removes all stored vectors and metadata from the Qdrant collection.
This operation is irreversible and will require re-ingestion of knowledge files.
Use with caution - this will make the knowledge base empty and RAG queries
will return no results until knowledge is re-ingested.
Returns:
Dict: Clear operation results including:
- Success status
- Operation message
- Collection cleared status
- Error details if applicable
Example Response:
```json
{
"success": true,
"message": "Knowledge base cleared successfully",
"collection_cleared": true
}
```
"""
try:
result = await knowledge_ingestion.clear_knowledge_base()
return result
except Exception as e:
return {
"success": False,
"error": str(e)
}
# Health check endpoint for monitoring
@app.get("/health", tags=["System Monitoring"])
async def health_check():
"""
System health check endpoint.
Provides basic system health information for monitoring and load balancers.
Returns:
Dict: Health status information
"""
return {
"status": "healthy",
"service": "Test the RAG System",
"version": "1.0.0",
"timestamp": time.time()
}
# Application startup and shutdown events
@app.on_event("startup")
async def startup_event():
"""Initialize system components on startup."""
print("π Test the RAG System starting up...")
print("π Knowledge processing system initialized")
print("π RAG pipeline ready")
print("π§ͺ RAG testing platform ready")
print("β
System ready for requests")
# Conversation Management Endpoints
@app.get("/conversation/status", tags=["Conversation Management"])
async def get_conversation_status():
"""
Get conversation tracking status and statistics.
Returns:
Dict: Conversation status information
"""
logger.info(f"π [API] /conversation/status endpoint called")
try:
stats = conversation_manager.get_conversation_stats()
logger.info(f"β
[CONVERSATION] Status retrieved successfully")
logger.info(f"π [CONVERSATION] Enabled: {stats['conversation_enabled']}")
logger.info(f"π [CONVERSATION] Total: {stats['total_conversations']}")
logger.info(f"π [CONVERSATION] Max: {stats['max_conversations']}")
logger.info(f"π [CONVERSATION] Serial range: {stats['oldest_serial']}-{stats['newest_serial']}")
return {
"conversation_enabled": stats["conversation_enabled"],
"total_conversations": stats["total_conversations"],
"max_conversations": stats["max_conversations"],
"oldest_serial": stats["oldest_serial"],
"newest_serial": stats["newest_serial"]
}
except Exception as e:
logger.error(f"β [API] /conversation/status failed: {str(e)}")
raise HTTPException(status_code=500, detail=f"Failed to get conversation status: {str(e)}")
@app.get("/conversation/history", tags=["Conversation Management"])
async def get_conversation_history():
"""
Get all stored conversation history.
Returns:
Dict: Complete conversation history
"""
conversations = conversation_manager.get_all_conversations()
return {
"conversations": conversations,
"total_count": len(conversations)
}
@app.delete("/conversation/clear", tags=["Conversation Management"])
async def clear_conversations():
"""
Clear all stored conversations.
Returns:
Dict: Confirmation message
"""
conversation_manager.clear_conversations()
return {
"message": "All conversations cleared successfully",
"success": True
}
@app.get("/conversation/config", tags=["Conversation Management"])
async def get_conversation_config():
"""
Get conversation configuration settings.
Returns:
Dict: Current conversation configuration
"""
config = get_config()
return {
"conversation_mood": config.conversation.conversation_mood,
"number_of_conversations_to_store": config.conversation.number_of_conversations_to_store,
"conversation_enabled": conversation_manager.is_conversation_enabled()
}
@app.on_event("shutdown")
async def shutdown_event():
"""Cleanup on shutdown."""
print("π Test the RAG System shutting down...")
print("β
Cleanup completed")
if __name__ == "__main__":
import uvicorn
uvicorn.run(
app,
host="0.0.0.0",
port=8000,
log_level="info",
access_log=True
)