Refactor code structure for improved readability and maintainability
This commit is contained in:
@@ -2,7 +2,7 @@
|
||||
|
||||
from fastapi import APIRouter
|
||||
|
||||
from app.api.v1 import auth, main
|
||||
from app.api.v1 import auth, main, socket
|
||||
|
||||
# V1 API router with v1 prefix
|
||||
api_router = APIRouter(prefix="/v1")
|
||||
@@ -10,3 +10,4 @@ api_router = APIRouter(prefix="/v1")
|
||||
# Include all route modules
|
||||
api_router.include_router(main.router, tags=["main"])
|
||||
api_router.include_router(auth.router, prefix="/auth", tags=["authentication"])
|
||||
api_router.include_router(socket.router, tags=["socket"])
|
||||
|
||||
66
app/api/v1/socket.py
Normal file
66
app/api/v1/socket.py
Normal file
@@ -0,0 +1,66 @@
|
||||
"""Socket.IO API endpoints for WebSocket management."""
|
||||
|
||||
from fastapi import APIRouter, Depends
|
||||
|
||||
from app.core.dependencies import get_current_user
|
||||
from app.models.user import User
|
||||
from app.services.socket import socket_manager
|
||||
|
||||
router = APIRouter(prefix="/socket", tags=["socket"])
|
||||
|
||||
|
||||
@router.get("/status")
|
||||
async def get_socket_status(current_user: User = Depends(get_current_user)):
|
||||
"""Get current socket connection status."""
|
||||
connected_users = socket_manager.get_connected_users()
|
||||
|
||||
return {
|
||||
"connected": str(current_user.id) in connected_users,
|
||||
"user_id": current_user.id,
|
||||
"total_connected": len(connected_users),
|
||||
}
|
||||
|
||||
|
||||
@router.post("/send-message")
|
||||
async def send_message_to_user(
|
||||
target_user_id: int,
|
||||
message: str,
|
||||
current_user: User = Depends(get_current_user),
|
||||
):
|
||||
"""Send a message to a specific user via WebSocket."""
|
||||
success = await socket_manager.send_to_user(
|
||||
str(target_user_id),
|
||||
"user_message",
|
||||
{
|
||||
"from_user_id": current_user.id,
|
||||
"from_user_name": current_user.name,
|
||||
"message": message,
|
||||
},
|
||||
)
|
||||
|
||||
return {
|
||||
"success": success,
|
||||
"target_user_id": target_user_id,
|
||||
"message": "Message sent" if success else "User not connected",
|
||||
}
|
||||
|
||||
|
||||
@router.post("/broadcast")
|
||||
async def broadcast_message(
|
||||
message: str,
|
||||
current_user: User = Depends(get_current_user),
|
||||
):
|
||||
"""Broadcast a message to all connected users."""
|
||||
await socket_manager.broadcast_to_all(
|
||||
"broadcast_message",
|
||||
{
|
||||
"from_user_id": current_user.id,
|
||||
"from_user_name": current_user.name,
|
||||
"message": message,
|
||||
},
|
||||
)
|
||||
|
||||
return {
|
||||
"success": True,
|
||||
"message": "Message broadcasted to all users",
|
||||
}
|
||||
@@ -1,6 +1,7 @@
|
||||
from collections.abc import AsyncGenerator
|
||||
from contextlib import asynccontextmanager
|
||||
|
||||
import socketio
|
||||
from fastapi import FastAPI
|
||||
from fastapi.middleware.cors import CORSMiddleware
|
||||
|
||||
@@ -8,6 +9,7 @@ from app.api import api_router
|
||||
from app.core.database import init_db
|
||||
from app.core.logging import get_logger, setup_logging
|
||||
from app.middleware.logging import LoggingMiddleware
|
||||
from app.services.socket import socket_manager
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
@@ -25,7 +27,7 @@ async def lifespan(_app: FastAPI) -> AsyncGenerator[None, None]:
|
||||
logger.info("Shutting down application")
|
||||
|
||||
|
||||
def create_app() -> FastAPI:
|
||||
def create_app():
|
||||
"""Create and configure the FastAPI application."""
|
||||
app = FastAPI(lifespan=lifespan)
|
||||
|
||||
@@ -43,7 +45,8 @@ def create_app() -> FastAPI:
|
||||
# Include API routes
|
||||
app.include_router(api_router)
|
||||
|
||||
return app
|
||||
# Create Socket.IO app with fallback to FastAPI app
|
||||
return socketio.ASGIApp(socket_manager.sio, app)
|
||||
|
||||
|
||||
app = create_app()
|
||||
|
||||
129
app/services/socket.py
Normal file
129
app/services/socket.py
Normal file
@@ -0,0 +1,129 @@
|
||||
"""WebSocket service for real-time communication with user rooms."""
|
||||
|
||||
import logging
|
||||
|
||||
import socketio
|
||||
|
||||
from app.utils.auth import JWTUtils
|
||||
from app.utils.cookies import extract_access_token_from_cookies
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class SocketManager:
|
||||
"""Manages WebSocket connections and user rooms."""
|
||||
|
||||
def __init__(self):
|
||||
self.sio = socketio.AsyncServer(
|
||||
cors_allowed_origins=["http://localhost:8001"],
|
||||
logger=True,
|
||||
engineio_logger=True,
|
||||
async_mode="asgi",
|
||||
)
|
||||
# Track user rooms: user_id -> room_id
|
||||
self.user_rooms: dict[str, str] = {}
|
||||
# Track socket sessions: socket_id -> user_id
|
||||
self.socket_users: dict[str, str] = {}
|
||||
|
||||
self._setup_handlers()
|
||||
|
||||
def _setup_handlers(self):
|
||||
"""Set up socket event handlers."""
|
||||
|
||||
@self.sio.event
|
||||
async def connect(sid, environ, auth=None):
|
||||
"""Handle client connection."""
|
||||
logger.info(f"Client {sid} attempting to connect")
|
||||
|
||||
# Extract access token from cookies
|
||||
cookie_header = environ.get("HTTP_COOKIE", "")
|
||||
access_token = extract_access_token_from_cookies(cookie_header)
|
||||
|
||||
if not access_token:
|
||||
logger.warning(f"Client {sid} connecting without access token")
|
||||
await self.sio.disconnect(sid)
|
||||
return
|
||||
|
||||
try:
|
||||
# Validate JWT token and extract user info
|
||||
payload = JWTUtils.decode_access_token(access_token)
|
||||
user_id = payload.get("sub")
|
||||
|
||||
if not user_id:
|
||||
logger.warning(f"Client {sid} token missing user ID")
|
||||
await self.sio.disconnect(sid)
|
||||
return
|
||||
|
||||
logger.info(f"User {user_id} connected with socket {sid}")
|
||||
except Exception as e:
|
||||
logger.warning(f"Client {sid} invalid token: {e}")
|
||||
await self.sio.disconnect(sid)
|
||||
return
|
||||
|
||||
# Store socket-user mapping
|
||||
self.socket_users[sid] = user_id
|
||||
|
||||
# Create or join user's personal room
|
||||
room_id = f"user_{user_id}"
|
||||
await self.sio.enter_room(sid, room_id)
|
||||
|
||||
# Update room tracking
|
||||
self.user_rooms[user_id] = room_id
|
||||
|
||||
logger.info(f"User {user_id} joined room {room_id}")
|
||||
|
||||
# Send welcome message to user
|
||||
await self.sio.emit(
|
||||
"user_connected",
|
||||
{
|
||||
"user_id": user_id,
|
||||
"room_id": room_id,
|
||||
"message": "Successfully connected to your personal room",
|
||||
},
|
||||
room=sid,
|
||||
)
|
||||
|
||||
@self.sio.event
|
||||
async def disconnect(sid):
|
||||
"""Handle client disconnection."""
|
||||
user_id = self.socket_users.get(sid)
|
||||
|
||||
if user_id:
|
||||
logger.info(f"User {user_id} disconnected (socket {sid})")
|
||||
# Clean up mappings
|
||||
del self.socket_users[sid]
|
||||
if user_id in self.user_rooms:
|
||||
del self.user_rooms[user_id]
|
||||
else:
|
||||
logger.info(f"Unknown client {sid} disconnected")
|
||||
|
||||
async def send_to_user(self, user_id: str, event: str, data: dict):
|
||||
"""Send a message to a specific user's room."""
|
||||
room_id = self.user_rooms.get(user_id)
|
||||
if room_id:
|
||||
await self.sio.emit(event, data, room=room_id)
|
||||
logger.debug(f"Sent {event} to user {user_id} in room {room_id}")
|
||||
return True
|
||||
else:
|
||||
logger.warning(f"User {user_id} not found in any room")
|
||||
return False
|
||||
|
||||
async def broadcast_to_all(self, event: str, data: dict):
|
||||
"""Broadcast a message to all connected users."""
|
||||
await self.sio.emit(event, data)
|
||||
logger.info(f"Broadcasted {event} to all users")
|
||||
|
||||
def get_connected_users(self) -> list:
|
||||
"""Get list of currently connected user IDs."""
|
||||
return list(self.user_rooms.keys())
|
||||
|
||||
def get_room_info(self) -> dict:
|
||||
"""Get information about connected users."""
|
||||
return {
|
||||
"total_users": len(self.user_rooms),
|
||||
"connected_users": list(self.user_rooms.keys()),
|
||||
}
|
||||
|
||||
|
||||
# Global socket manager instance
|
||||
socket_manager = SocketManager()
|
||||
24
app/utils/cookies.py
Normal file
24
app/utils/cookies.py
Normal file
@@ -0,0 +1,24 @@
|
||||
"""Cookie parsing utilities for WebSocket authentication."""
|
||||
|
||||
from typing import Optional
|
||||
|
||||
|
||||
def parse_cookies(cookie_header: str) -> dict[str, str]:
|
||||
"""Parse HTTP cookie header into a dictionary."""
|
||||
cookies = {}
|
||||
if not cookie_header:
|
||||
return cookies
|
||||
|
||||
for cookie in cookie_header.split(";"):
|
||||
cookie = cookie.strip()
|
||||
if "=" in cookie:
|
||||
name, value = cookie.split("=", 1)
|
||||
cookies[name.strip()] = value.strip()
|
||||
|
||||
return cookies
|
||||
|
||||
|
||||
def extract_access_token_from_cookies(cookie_header: str) -> Optional[str]:
|
||||
"""Extract access token from HTTP cookies."""
|
||||
cookies = parse_cookies(cookie_header)
|
||||
return cookies.get("access_token")
|
||||
Reference in New Issue
Block a user