Spaces:
Sleeping
Sleeping
Ekimart commited on
Commit ·
527670e
1
Parent(s): f2a407c
Mettre à jour WebSocketManager pour implémenter une structure de message améliorée, gérer les pings/pongs, et ajouter des ajustements de configuration pour les connexions WebSocket.
Browse files- core/manager/websocket_manager.py +8 -7
- main.py +21 -23
core/manager/websocket_manager.py
CHANGED
|
@@ -34,8 +34,9 @@ class WebSocketManager:
|
|
| 34 |
await self._send(
|
| 35 |
websocket,
|
| 36 |
{
|
| 37 |
-
"type": "connection_established",
|
| 38 |
"session_id": session_id,
|
|
|
|
|
|
|
| 39 |
"message": "Connexion au service de notifications établie",
|
| 40 |
"timestamp": current_utc_iso(),
|
| 41 |
},
|
|
@@ -52,7 +53,7 @@ class WebSocketManager:
|
|
| 52 |
self.logger.info(f"Connexion fermée: {session_id}")
|
| 53 |
|
| 54 |
async def send_personal_message(
|
| 55 |
-
|
| 56 |
) -> None:
|
| 57 |
"""Envoie un message à un client spécifique"""
|
| 58 |
if ws := self.active_connections.get(session_id):
|
|
@@ -152,8 +153,8 @@ class WebSocketManager:
|
|
| 152 |
"timestamp": current_utc_iso(),
|
| 153 |
}
|
| 154 |
|
| 155 |
-
|
| 156 |
-
async def _send(
|
| 157 |
"""Mécanisme d'envoi générique"""
|
| 158 |
await websocket.send_json(payload)
|
| 159 |
|
|
@@ -164,9 +165,9 @@ class WebSocketManager:
|
|
| 164 |
|
| 165 |
@staticmethod
|
| 166 |
def _should_receive(
|
| 167 |
-
|
| 168 |
-
|
| 169 |
-
|
| 170 |
) -> bool:
|
| 171 |
"""Vérifie si un client doit recevoir une notification"""
|
| 172 |
return any(
|
|
|
|
| 34 |
await self._send(
|
| 35 |
websocket,
|
| 36 |
{
|
|
|
|
| 37 |
"session_id": session_id,
|
| 38 |
+
"type": "connection_established",
|
| 39 |
+
"title": "Connexion réussie",
|
| 40 |
"message": "Connexion au service de notifications établie",
|
| 41 |
"timestamp": current_utc_iso(),
|
| 42 |
},
|
|
|
|
| 53 |
self.logger.info(f"Connexion fermée: {session_id}")
|
| 54 |
|
| 55 |
async def send_personal_message(
|
| 56 |
+
self, session_id: str, payload: Dict[str, Any]
|
| 57 |
) -> None:
|
| 58 |
"""Envoie un message à un client spécifique"""
|
| 59 |
if ws := self.active_connections.get(session_id):
|
|
|
|
| 153 |
"timestamp": current_utc_iso(),
|
| 154 |
}
|
| 155 |
|
| 156 |
+
@staticmethod
|
| 157 |
+
async def _send(websocket: WebSocket, payload: Dict[str, Any]) -> None:
|
| 158 |
"""Mécanisme d'envoi générique"""
|
| 159 |
await websocket.send_json(payload)
|
| 160 |
|
|
|
|
| 165 |
|
| 166 |
@staticmethod
|
| 167 |
def _should_receive(
|
| 168 |
+
subscriptions: List[Subscription],
|
| 169 |
+
notification: Notification,
|
| 170 |
+
client: ClientInfo,
|
| 171 |
) -> bool:
|
| 172 |
"""Vérifie si un client doit recevoir une notification"""
|
| 173 |
return any(
|
main.py
CHANGED
|
@@ -13,6 +13,7 @@ from fastapi import (
|
|
| 13 |
Depends,
|
| 14 |
)
|
| 15 |
from fastapi.templating import Jinja2Templates
|
|
|
|
| 16 |
|
| 17 |
from core.manager.websocket_manager import WebSocketManager
|
| 18 |
from core.settings.settings import Settings
|
|
@@ -61,6 +62,8 @@ app = FastAPI(
|
|
| 61 |
lifespan=lifespan,
|
| 62 |
docs_url="/docs",
|
| 63 |
redoc_url="/redoc",
|
|
|
|
|
|
|
| 64 |
)
|
| 65 |
|
| 66 |
|
|
@@ -83,16 +86,11 @@ async def websocket_endpoint(websocket: WebSocket, user_id: str):
|
|
| 83 |
session_id = await manager.connect(websocket, client_info)
|
| 84 |
|
| 85 |
while True:
|
| 86 |
-
|
| 87 |
-
|
| 88 |
-
|
| 89 |
-
|
| 90 |
-
|
| 91 |
-
"type": "message_ack",
|
| 92 |
-
"message": f"Message reçu: {data}",
|
| 93 |
-
"timestamp": current_utc_iso(),
|
| 94 |
-
},
|
| 95 |
-
)
|
| 96 |
|
| 97 |
except WebSocketDisconnect:
|
| 98 |
manager.disconnect(session_id)
|
|
@@ -103,24 +101,24 @@ async def websocket_endpoint(websocket: WebSocket, user_id: str):
|
|
| 103 |
|
| 104 |
|
| 105 |
# Endpoints REST
|
| 106 |
-
@app.get("/")
|
| 107 |
async def root(request: Request):
|
| 108 |
-
|
| 109 |
-
|
| 110 |
-
|
| 111 |
-
|
| 112 |
-
|
| 113 |
-
return manager.get_stats()
|
| 114 |
|
| 115 |
|
| 116 |
@app.post("/subscriptions/{session_id}")
|
| 117 |
async def add_subscription(
|
| 118 |
-
|
| 119 |
-
|
| 120 |
):
|
| 121 |
"""Ajoute un nouvel abonnement pour une session"""
|
| 122 |
|
| 123 |
-
subscription =
|
|
|
|
| 124 |
|
| 125 |
if not manager.add_subscription(session_id, subscription):
|
| 126 |
raise HTTPException(status_code=400, detail="Échec de l'ajout de l'abonnement")
|
|
@@ -133,8 +131,8 @@ async def add_subscription(
|
|
| 133 |
|
| 134 |
@app.delete("/subscriptions/{session_id}/{subscription_id}")
|
| 135 |
async def remove_subscription(
|
| 136 |
-
|
| 137 |
-
|
| 138 |
):
|
| 139 |
"""Supprime un abonnement existant"""
|
| 140 |
if not manager.remove_subscription(session_id, subscription_id):
|
|
@@ -151,7 +149,7 @@ async def get_subscriptions(session_id: str = Depends(validate_session)):
|
|
| 151 |
|
| 152 |
@app.post("/notifications/user/{user_id}")
|
| 153 |
async def send_user_notification(
|
| 154 |
-
|
| 155 |
):
|
| 156 |
"""Envoie une notification à un utilisateur spécifique"""
|
| 157 |
notification = Notification(**notification_request.model_dump())
|
|
|
|
| 13 |
Depends,
|
| 14 |
)
|
| 15 |
from fastapi.templating import Jinja2Templates
|
| 16 |
+
from starlette.responses import HTMLResponse
|
| 17 |
|
| 18 |
from core.manager.websocket_manager import WebSocketManager
|
| 19 |
from core.settings.settings import Settings
|
|
|
|
| 62 |
lifespan=lifespan,
|
| 63 |
docs_url="/docs",
|
| 64 |
redoc_url="/redoc",
|
| 65 |
+
websocket_ping_interval=20,
|
| 66 |
+
websocket_timeout=60,
|
| 67 |
)
|
| 68 |
|
| 69 |
|
|
|
|
| 86 |
session_id = await manager.connect(websocket, client_info)
|
| 87 |
|
| 88 |
while True:
|
| 89 |
+
message = await websocket.receive_text()
|
| 90 |
+
if message == "ping":
|
| 91 |
+
await websocket.send_text("pong")
|
| 92 |
+
else:
|
| 93 |
+
pass
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 94 |
|
| 95 |
except WebSocketDisconnect:
|
| 96 |
manager.disconnect(session_id)
|
|
|
|
| 101 |
|
| 102 |
|
| 103 |
# Endpoints REST
|
| 104 |
+
@app.get("/", response_class=HTMLResponse)
|
| 105 |
async def root(request: Request):
|
| 106 |
+
"""Page de test pour les connexions WebSocket"""
|
| 107 |
+
return templates.TemplateResponse(
|
| 108 |
+
"dashboard.html",
|
| 109 |
+
{"request": request, "current_time": current_utc_iso()},
|
| 110 |
+
)
|
|
|
|
| 111 |
|
| 112 |
|
| 113 |
@app.post("/subscriptions/{session_id}")
|
| 114 |
async def add_subscription(
|
| 115 |
+
session_id: str = Depends(validate_session),
|
| 116 |
+
subscription_request: Optional[SubscriptionRequest] = None,
|
| 117 |
):
|
| 118 |
"""Ajoute un nouvel abonnement pour une session"""
|
| 119 |
|
| 120 |
+
subscription = Subscription(**subscription_request.model_dump()) if subscription_request else Subscription(
|
| 121 |
+
type="default", target="default")
|
| 122 |
|
| 123 |
if not manager.add_subscription(session_id, subscription):
|
| 124 |
raise HTTPException(status_code=400, detail="Échec de l'ajout de l'abonnement")
|
|
|
|
| 131 |
|
| 132 |
@app.delete("/subscriptions/{session_id}/{subscription_id}")
|
| 133 |
async def remove_subscription(
|
| 134 |
+
subscription_id: str,
|
| 135 |
+
session_id: str = Depends(validate_session),
|
| 136 |
):
|
| 137 |
"""Supprime un abonnement existant"""
|
| 138 |
if not manager.remove_subscription(session_id, subscription_id):
|
|
|
|
| 149 |
|
| 150 |
@app.post("/notifications/user/{user_id}")
|
| 151 |
async def send_user_notification(
|
| 152 |
+
user_id: str, notification_request: NotificationRequest
|
| 153 |
):
|
| 154 |
"""Envoie une notification à un utilisateur spécifique"""
|
| 155 |
notification = Notification(**notification_request.model_dump())
|