Serveur et Clients WebSocket pour la Messagerie Instantanée

La communication en temps réel est devenue un pilier essentiel des applications web modernes, facilitant des interactions dynamiques telles que les chats en ligne, les notifications push instantanées ou la mise à jour en direct des cotations boursières. Cet article détaille la mise en œuvre d'un système de diffusion de messages instantanés, robuste et performant, en s'appuyant sur la technologie WebSocket.

Choix Technologiques Clés :

  • Serveur : Go, exploitant la bibliothèque gorilla/websocket pour la gestion des connexions.
  • Clients : Un exemple de client en Python (utilisant la bibliothèque websockets) et un autre en Go.
  • Protocole : WebSocket pour une communication bidirectionnelle persistante.
  • Gestion de la Concurrence : Utilisation de sync.Map en Go pour maintenir une collection thread-safe des clients connectés, garantissant une haute scalabilité.
  • Fiabilité : Intégration de mécanismes de détection de déconnexion, de gestion des timeouts et de signaux de maintien en vie (heartbeats).

Implémentation du Serveur WebSocket en Go

Le serveur écoute sur l'adresse :8808/chat. Il est conçu pour gérer de multiples connexions simultanées, diffuser des messages à tous les clients connectés et assurer la robustesse des liens grâce à un mécanisme de 'heartbeat'.

package main

import (
	"context"
	"log"
	"net/http"
	"sync"
	"time"

	"github.com/gorilla/websocket"
)

// ClientSession représente une session WebSocket client, sécurisant les écritures.
type ClientSession struct {
	Conn *websocket.Conn
	mu   sync.Mutex // Mutex pour protéger les opérations d'écriture sur la connexion.
}

var (
	// connectedClients stocke toutes les sessions client actives.
	connectedClients sync.Map
	// messageRelay est un canal pour la diffusion de messages à tous les clients.
	messageRelay = make(chan []byte)
)

// wsUpgrader est la configuration pour la mise à niveau des requêtes HTTP en connexions WebSocket.
var wsUpgrader = websocket.Upgrader{
	CheckOrigin: func(r *http.Request) bool {
		// Autoriser toutes les origines pour simplifier l'exemple. En production, cela devrait être plus restrictif.
		return true
	},
	ReadBufferSize:  1024,
	WriteBufferSize: 1024,
}

func main() {
	// Créer un contexte pour gérer la terminaison du programme.
	ctx, cancel := context.WithCancel(context.Background())
	defer cancel()

	http.HandleFunc("/chat", handleWebSocketConnection)
	go processAndRelayMessages(ctx) // Démarrer la goroutine de diffusion des messages.

	log.Println("Serveur WebSocket démarré sur localhost:8808")
	if err := http.ListenAndServe(":8808", nil); err != nil {
		log.Fatalf("Échec du démarrage du serveur HTTP: %v", err)
	}
}

// handleWebSocketConnection gère les nouvelles connexions WebSocket.
func handleWebSocketConnection(w http.ResponseWriter, r *http.Request) {
	conn, err := wsUpgrader.Upgrade(w, r, nil)
	if err != nil {
		log.Printf("Erreur lors de la mise à niveau de la connexion: %v", err)
		return
	}

	session := &ClientSession{Conn: conn}
	connID := conn.RemoteAddr().String()
	connectedClients.Store(connID, session)
	log.Printf("Nouvelle connexion établie : %s", connID)

	defer func() {
		connectedClients.Delete(connID)
		conn.Close()
		log.Printf("Connexion terminée pour : %s", connID)
	}()

	// Configurer la gestion des Pongs pour maintenir la connexion active.
	// La date limite de lecture est étendue chaque fois qu'un Pong est reçu.
	conn.SetReadDeadline(time.Now().Add(60 * time.Second))
	conn.SetPongHandler(func(string) error {
		conn.SetReadDeadline(time.Now().Add(60 * time.Second))
		return nil
	})

	// Lancer une goroutine pour envoyer des pings de maintien en vie.
	go func() {
		pingTicker := time.NewTicker(30 * time.Second)
		defer pingTicker.Stop()
		for {
			select {
			case <-pingTicker.C:
				session.mu.Lock()
				err := session.Conn.WriteMessage(websocket.PingMessage, nil)
				session.mu.Unlock()
				if err != nil {
					log.Printf("Erreur lors de l'envoi du Ping à %s: %v", connID, err)
					// Si le ping échoue, considérer la connexion comme rompue.
					conn.Close()
					connectedClients.Delete(connID)
					return // Terminer cette goroutine.
				}
			}
		}
	}()

	// Boucle principale de lecture des messages du client.
	for {
		_, payload, err := conn.ReadMessage()
		if err != nil {
			// Gérer les erreurs de lecture (ex: client déconnecté).
			if websocket.IsCloseError(err, websocket.CloseGoingAway, websocket.CloseNormalClosure) {
				log.Printf("Client %s a fermé la connexion normalement.", connID)
			} else {
				log.Printf("Erreur de lecture du message pour %s: %v", connID, err)
			}
			break // Sortir de la boucle de lecture.
		}
		log.Printf("Message reçu de %s: %s", connID, string(payload))
		messageRelay <- payload // Envoyer le message au canal de diffusion.
	}
}

// processAndRelayMessages gère la distribution des messages reçus à tous les clients.
func processAndRelayMessages(ctx context.Context) {
	for {
		select {
		case payload := <-messageRelay:
			distributeToAllClients(payload)
		case <-ctx.Done():
			log.Println("Arrêt de la routine de diffusion des messages.")
			return
		}
	}
}

// distributeToAllClients envoie le message donné à tous les clients connectés.
func distributeToAllClients(payload []byte) {
	connectedClients.Range(func(key, value interface{}) bool {
		session := value.(*ClientSession)
		session.mu.Lock()
		defer session.mu.Unlock()

		if err := session.Conn.WriteMessage(websocket.TextMessage, payload); err != nil {
			log.Printf("Erreur d'écriture message à %s: %v", key, err)
			// Si l'écriture échoue, la connexion est probablement morte.
			session.Conn.Close()
			connectedClients.Delete(key)
		}
		return true // Continuer à itérer sur les clients.
	})
}

Client WebSocket en Python (Émetteur)

Ce client Python est conçu pour envoyer des messages structurés au serveur WebSocket et, si applicable, attendre une réponse dans un délai spécifié.

import asyncio
import json
import logging
import datetime
import websockets
import time # Nécessaire pour time.sleep dans le bloc principal
from typing import Dict, Any

# Configuration du logger pour une meilleure visibilité.
_logger = logging.getLogger(__name__)
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(levelname)s - %(message)s')

async def transmit_ws_message(endpoint_url: str, data_payload: Dict[str, Any], response_timeout_sec: int = 10) -> Dict[str, Any]:
    """
    Établit une connexion WebSocket, envoie un message JSON et attend une réponse.

    Args:
        endpoint_url (str): L'URL du serveur WebSocket (ex: "ws://127.0.0.1:8808/chat").
        data_payload (Dict[str, Any]): Le dictionnaire de données à envoyer, qui sera sérialisé en JSON.
        response_timeout_sec (int): Le temps maximal en secondes pour attendre une réponse.

    Returns:
        Dict[str, Any]: La réponse décodée du serveur ou un dictionnaire d'erreur.
    """
    try:
        async with websockets.connect(endpoint_url) as ws_conn:
            await asyncio.wait_for(ws_conn.send(json.dumps(data_payload)), response_timeout_sec)
            _logger.info(f"Message envoyé: {data_payload}")

            # Le serveur actuel diffuse tous les messages à tous les clients,
            # donc ce client peut recevoir son propre message ou ceux d'autres clients.
            # Nous attendons ici une réponse comme si le serveur renvoyait un accusé de réception.
            server_response_raw = await asyncio.wait_for(ws_conn.recv(), response_timeout_sec)
            server_response = json.loads(server_response_raw)
            return server_response

    except asyncio.TimeoutError:
        _logger.error(f"Délai d'attente dépassé après {response_timeout_sec}s pour l'URL {endpoint_url}")
        return {"error": "Connection or response timed out"}
    except websockets.exceptions.ConnectionClosedOK:
        _logger.warning("Connexion WebSocket fermée normalement par le serveur.")
        return {"info": "Connection closed normally"}
    except websockets.exceptions.ConnectionClosedError as e:
        _logger.error(f"Connexion WebSocket fermée de manière inattendue: {e}")
        return {"error": f"Connection closed unexpectedly: {e}"}
    except json.JSONDecodeError as e:
        _logger.error(f"Erreur de décodage JSON de la réponse du serveur: {e}")
        return {"error": f"JSON decode error in server response: {e}"}
    except Exception as e:
        _logger.critical(f"Erreur inattendue lors de la transmission du message: {e}")
        return {"error": f"An unexpected error occurred: {e}"}

def send_message_blocking(payload_data: Dict[str, Any]) -> Dict[str, Any]:
    """
    Fonction bloquante pour envoyer un message WebSocket.
    Utilise asyncio.run pour exécuter la coroutine asynchrone.
    """
    return asyncio.run(transmit_ws_message("ws://127.0.0.1:8808/chat", payload_data))

if __name__ == "__main__":
    _logger.info("Démarrage du client WebSocket Python (émetteur).")
    for i in range(5): # Réduit le nombre d'itérations pour un test rapide.
        sample_payload = {
            "origin": "PythonClientA",
            "sequence": i + 1,
            "timestamp": datetime.datetime.now().isoformat(),
            "content": f"Bonjour du client Python, message numéro {i+1}!"
        }
        response = send_message_blocking(sample_payload)
        _logger.info(f"Réponse ou statut de l'envoi: {response}")
        time.sleep(1) # Pause d'une seconde entre les envois.

Client WebSocket en Go (Récepteur)

Ce client Go est configuré pour se connecter au serveur WebSocket et écouter en continu les messages diffusés.

package main

import (
	"fmt"
	"log"
	"os"
	"os/signal"
	"syscall"
	"time"

	"github.com/gorilla/websocket"
)

func main() {
	serverURL := "ws://127.0.0.1:8808/chat"
	log.Printf("Tentative de connexion au serveur WebSocket : %s", serverURL)

	// Établir la connexion WebSocket
	wsConn, _, err := websocket.DefaultDialer.Dial(serverURL, nil)
	if err != nil {
		log.Fatalf("Échec de la connexion au serveur %s: %v", serverURL, err)
	}
	defer wsConn.Close() // S'assurer que la connexion est fermée à la fin du programme.

	// Configuration pour intercepter les signaux d'interruption du système (Ctrl+C, etc.)
	signalChannel := make(chan os.Signal, 1)
	signal.Notify(signalChannel, syscall.SIGINT, syscall.SIGTERM)

	// Lancer une goroutine pour la lecture asynchrone des messages.
	go func() {
		log.Println("Commence l'écoute des messages du serveur...")
		for {
			messageType, receivedData, err := wsConn.ReadMessage()
			if err != nil {
				// Gérer les erreurs de lecture, y compris les fermetures de connexion.
				if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseNormalClosure) {
					log.Printf("Erreur de lecture inattendue : %v", err)
				} else {
					log.Printf("Connexion fermée ou erreur de lecture : %v", err)
				}
				break // Sortir de la boucle si la connexion est fermée ou une erreur irrécupérable survient.
			}
			// Afficher le contenu du message en fonction de son type.
			switch messageType {
			case websocket.TextMessage:
				fmt.Printf("Message reçu du serveur : %s\n", string(receivedData))
			case websocket.BinaryMessage:
				fmt.Printf("Message binaire reçu (taille %d octets)\n", len(receivedData))
			case websocket.CloseMessage:
				fmt.Println("Message de fermeture de connexion reçu du serveur.")
				return // Terminer la goroutine de lecture.
			default:
				fmt.Printf("Type de message inconnu reçu (%d), données : %s\n", messageType, string(receivedData))
			}
		}
	}()

	// Le programme principal attend un signal d'interruption.
	<-signalChannel
	log.Println("Signal d'interruption reçu. Tentative de fermeture de la connexion WebSocket...")

	// Envoyer un message de fermeture propre au serveur.
	err = wsConn.WriteMessage(websocket.CloseMessage, websocket.FormatCloseMessage(websocket.CloseNormalClosure, ""))
	if err != nil {
		log.Printf("Erreur lors de l'envoi du message de fermeture: %v", err)
	}

	// Donner un court délai pour que le message de fermeture soit traité.
	time.Sleep(time.Second)
	log.Println("Client WebSocket Go arrêté.")
}

Synthèse des Fnoctionnalités et Cas d'Usage

Ce projet démontre la construction d'un système de messagerie en temps réel basé sur WebSocket, offrant les capacités clés suivantes :

  • Gestion multi-clients : Prise en charge de nombreuses connexions simultanées.
  • Diffusion de messages : Envoi efficace d'informations à l'ensemble des clients connectés.
  • Maintien en vie (Heartbeat) : Mécanisme robuste pour détecter et gérer les déconnexions.
  • Gestion des erreurs : Traitement des incidents liés à la connexion ou à la transmission.

Une telle architecture WebSocket est parfaitement adaptée à divers scénarios exigeant une réactivité immédiate, notamment les applications de chat instantané, les systèmes de notifications push, les outils de collaboration en temps réel (édition collaborative), et les plateformes de suivi d'événements en direct.

Étiquettes: WebSocket Go Python real-time Concurrency

Publié le 4 septembre à 19h49