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
133 changes: 133 additions & 0 deletions app/api/monitoring.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,133 @@
"""Monitoring endpoints for system health and metrics."""

import time

from fastapi import APIRouter, Response

from ..services.cache.redis_manager import redis_manager
from ..services.cache.semantic_cache import semantic_cache
from ..services.embeddings.embedding_service import embedding_service
from ..utils.logger import logger
from ..utils.metrics import metrics_collector

router = APIRouter(prefix="/monitoring", tags=["Monitoring"])


@router.get("/cache/stats")
async def get_cache_statistics():
"""Get comprehensive cache performance statistics."""
try:
redis_info = await redis_manager.get_info()
semantic_stats = semantic_cache.get_stats()
embedding_stats = embedding_service.get_stats()

return {
"status": "healthy",
"redis": redis_info,
"semantic_cache": semantic_stats,
"embeddings": embedding_stats,
"recommendations": _get_cache_recommendations(semantic_stats),
}
except Exception as e:
logger.error(f"Failed to get cache statistics: {e}")
return {
"status": "error",
"error": str(e),
}


@router.get("/cache/health")
async def check_cache_health():
"""Quick health check for cache systems."""
try:
redis_healthy = await redis_manager.exists("health_check")

return {
"redis": "healthy" if redis_healthy or redis_manager._is_healthy else "unhealthy",
"status": "healthy" if redis_healthy or redis_manager._is_healthy else "degraded",
}
except Exception as e:
logger.error(f"Cache health check failed: {e}")
return {
"status": "unhealthy",
"error": str(e),
}


@router.delete("/cache/clear")
async def clear_cache():
"""Clear all cache entries (admin operation)."""
try:
semantic_cleared = await semantic_cache.clear_all()

embedding_cleared = await embedding_service.clear_cache()

logger.info(
"Cache cleared",
extra={
"semantic_entries": semantic_cleared,
"embedding_entries": embedding_cleared,
},
)

return {
"status": "success",
"semantic_entries_cleared": semantic_cleared,
"embedding_entries_cleared": embedding_cleared,
"total_cleared": semantic_cleared + embedding_cleared,
}
except Exception as e:
logger.error(f"Failed to clear cache: {e}")
return {
"status": "error",
"error": str(e),
}


def _get_cache_recommendations(stats: dict) -> list[str]:
"""Generate recommendations based on cache statistics."""
recommendations = []

if stats["hit_rate"] < 0.2:
recommendations.append(
"Low cache hit rate. Consider adjusting similarity threshold or warming cache with common patterns."
)

if stats["hit_rate"] > 0.9:
recommendations.append("Very high cache hit rate. Consider reducing TTL to ensure fresh responses.")

if stats["total_requests"] > 10000:
recommendations.append(
"High cache usage. Monitor memory consumption and consider implementing cache size limits."
)

return recommendations


@router.get("/metrics")
async def get_prometheus_metrics():
"""Get Prometheus metrics in text format."""
try:
metrics_data = metrics_collector.get_metrics()
return Response(content=metrics_data, media_type="text/plain")
except Exception as e:
logger.error(f"Failed to get Prometheus metrics: {e}")
return Response(content=f"# Error: {str(e)}", media_type="text/plain", status_code=500)


@router.get("/metrics/json")
async def get_metrics_json():
"""Get metrics in JSON format for easier consumption."""
try:
metrics_dict = metrics_collector.get_metrics_dict()
return {
"status": "success",
"metrics": metrics_dict,
"timestamp": int(time.time()),
}
except Exception as e:
logger.error(f"Failed to get metrics: {e}")
return {
"status": "error",
"error": str(e),
}
110 changes: 109 additions & 1 deletion app/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,9 +10,12 @@

from app.api import chat as chat_api
from app.api import conversation_analysis as analysis_api
from app.api import monitoring as monitoring_api
from app.api import user as user_api
from app.db.chat import db
from app.services.cache.redis_manager import redis_manager
from app.utils.logger import logger
from app.utils.metrics import metrics_collector


@asynccontextmanager
Expand All @@ -21,10 +24,21 @@ async def lifespan(app: FastAPI):
Handles startup and shutdown events for the application.
"""
logger.info("Starting up...")

await db.create_db_and_tables()
logger.info("Database tables created or already exist.")

await redis_manager.initialize()
logger.info("Redis cache initialized.")

metrics_collector.initialize()
logger.info("Metrics collector initialized.")

yield

logger.info("Shutting down...")
await redis_manager.close()
logger.info("Redis connections closed.")


app = FastAPI(
Expand Down Expand Up @@ -142,12 +156,106 @@ async def general_exception_handler(request: Request, exc: Exception):

@app.get("/health")
async def health_check():
return {"status": "healthy"}
"""Health check endpoint that validates all services."""
import time

from .services.cache.redis_manager import redis_manager
from .services.cache.semantic_cache import semantic_cache
from .utils.config import app_settings

health_status = {
"status": "healthy",
"timestamp": int(time.time()),
"version": app_settings.VERSION if hasattr(app_settings, "VERSION") else "unknown",
"services": {},
}

try:
redis_healthy = redis_manager.is_healthy
health_status["services"]["redis"] = {"status": "healthy" if redis_healthy else "unhealthy"}
if not redis_healthy:
health_status["status"] = "unhealthy"
except Exception as e:
health_status["services"]["redis"] = {"status": "unhealthy", "error": str(e)}
health_status["status"] = "unhealthy"

try:
cache_healthy = await semantic_cache.health_check()
health_status["services"]["cache"] = {"status": "healthy" if cache_healthy else "unhealthy"}
if not cache_healthy:
health_status["status"] = "unhealthy"
health_status["services"]["cache"]["error"] = "Cache health check failed"
except Exception as e:
health_status["services"]["cache"] = {"status": "unhealthy", "error": str(e)}
health_status["status"] = "unhealthy"

if health_status["status"] == "unhealthy":
return JSONResponse(status_code=503, content=health_status)

return health_status


@app.get("/health/detailed")
async def health_check_detailed():
"""Detailed health check with additional service information."""

from .services.cache.redis_manager import redis_manager
from .services.cache.semantic_cache import semantic_cache

health_status = await health_check()
if isinstance(health_status, JSONResponse):
health_status = health_status.body.decode()
import json

health_status = json.loads(health_status)

try:
if hasattr(redis_manager, "get_connection_info"):
health_status["services"]["redis"]["connection_info"] = redis_manager.get_connection_info()
except Exception:
pass

try:
if hasattr(semantic_cache, "get_stats"):
health_status["services"]["cache"]["stats"] = await semantic_cache.get_stats()
except Exception:
pass

return health_status


@app.get("/metrics")
async def get_metrics():
"""Prometheus metrics endpoint."""
from fastapi import Response

from .utils.metrics import metrics_collector

try:
metrics_data = metrics_collector.get_metrics()
return Response(content=metrics_data, media_type="text/plain; version=0.0.4")
except Exception as e:
logger.error(f"Failed to get metrics: {e}")
return Response(content=f"# Error: {str(e)}", media_type="text/plain; version=0.0.4", status_code=500)


@app.get("/metrics/json")
async def get_metrics_json():
"""JSON metrics endpoint."""
from .utils.metrics import metrics_collector

try:
metrics_dict = metrics_collector.get_metrics_dict()
return metrics_dict
except Exception as e:
logger.error(f"Failed to get metrics: {e}")
return JSONResponse(status_code=500, content={"error": str(e)})


app.include_router(user_api.router)
app.include_router(chat_api.router)
app.include_router(analysis_api.router)
app.include_router(monitoring_api.router)


def run_uvicorn():
Expand Down
40 changes: 40 additions & 0 deletions app/services/cache/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
"""Cache services for the Conversational Analysis Engine."""

from .cache_metrics import cache_metrics, track_cache_operation
from .eviction_policies import (
EvictionPolicy,
EvictionPolicyFactory,
HybridEvictionPolicy,
LFUEvictionPolicy,
LRUEvictionPolicy,
TTLEvictionPolicy,
)
from .redis_manager import redis_manager
from .semantic_cache import semantic_cache
from .similarity_strategies import (
CosineSimilarityStrategy,
DotProductSimilarityStrategy,
EuclideanDistanceStrategy,
HybridSimilarityStrategy,
SimilarityStrategy,
SimilarityStrategyFactory,
)

__all__ = [
"redis_manager",
"semantic_cache",
"cache_metrics",
"track_cache_operation",
"EvictionPolicy",
"TTLEvictionPolicy",
"LRUEvictionPolicy",
"LFUEvictionPolicy",
"HybridEvictionPolicy",
"EvictionPolicyFactory",
"SimilarityStrategy",
"CosineSimilarityStrategy",
"EuclideanDistanceStrategy",
"DotProductSimilarityStrategy",
"HybridSimilarityStrategy",
"SimilarityStrategyFactory",
]
Loading