🔄 WebSocket: Full-Duplex
WebSocket provides bidirectional, persistent connection. Best for chat, gaming, real-time collaboration.
Traditional HTTP: Client requests, server responds. Client must poll for updates.
Problems:
Solution: Real-time communication - Server pushes updates to client!
Bidirectional, persistent connection.
Characteristics:
Server-to-client streaming over HTTP.
Characteristics:
Hold request open until data available.
Characteristics:
1. Handshake (HTTP Upgrade):
GET /ws HTTP/1.1Host: example.comUpgrade: websocketConnection: UpgradeSec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ==Sec-WebSocket-Version: 13Server responds:
HTTP/1.1 101 Switching ProtocolsUpgrade: websocketConnection: UpgradeSec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=2. Connection Established - Now bidirectional!
import asyncioimport websocketsimport json
class WebSocketServer: def __init__(self): self.clients = set()
async def register_client(self, websocket): """Register new client""" self.clients.add(websocket) print(f"Client connected. Total: {len(self.clients)}")
async def unregister_client(self, websocket): """Unregister client""" self.clients.discard(websocket) print(f"Client disconnected. Total: {len(self.clients)}")
async def handle_client(self, websocket, path): """Handle client connection""" await self.register_client(websocket)
try: async for message in websocket: # Handle incoming message data = json.loads(message) await self.process_message(websocket, data)
except websockets.exceptions.ConnectionClosed: pass finally: await self.unregister_client(websocket)
async def process_message(self, websocket, data): """Process client message""" message_type = data.get('type')
if message_type == 'chat': # Broadcast to all clients await self.broadcast({ 'type': 'chat', 'user': data.get('user'), 'message': data.get('message'), 'timestamp': data.get('timestamp') })
elif message_type == 'ping': # Respond to ping await websocket.send(json.dumps({'type': 'pong'}))
async def broadcast(self, message): """Broadcast message to all clients""" if self.clients: message_json = json.dumps(message) await asyncio.gather( *[client.send(message_json) for client in self.clients], return_exceptions=True )
async def send_to_client(self, websocket, message): """Send message to specific client""" try: await websocket.send(json.dumps(message)) except websockets.exceptions.ConnectionClosed: await self.unregister_client(websocket)
# Start serverserver = WebSocketServer()start_server = websockets.serve( server.handle_client, "localhost", 8765)
asyncio.get_event_loop().run_until_complete(start_server)asyncio.get_event_loop().run_forever()import org.java_websocket.WebSocket;import org.java_websocket.handshake.ClientHandshake;import org.java_websocket.server.WebSocketServer;import java.net.InetSocketAddress;import java.util.Collections;import java.util.Set;import java.util.concurrent.CopyOnWriteArraySet;
public class ChatServer extends WebSocketServer { private final Set<WebSocket> clients = new CopyOnWriteArraySet<>();
public ChatServer(int port) { super(new InetSocketAddress(port)); }
@Override public void onOpen(WebSocket conn, ClientHandshake handshake) { // Register new client clients.add(conn); System.out.println("Client connected. Total: " + clients.size()); }
@Override public void onClose(WebSocket conn, int code, String reason, boolean remote) { // Unregister client clients.remove(conn); System.out.println("Client disconnected. Total: " + clients.size()); }
@Override public void onMessage(WebSocket conn, String message) { // Handle incoming message try { JSONObject data = new JSONObject(message); String type = data.getString("type");
if ("chat".equals(type)) { // Broadcast to all clients broadcast(message); } else if ("ping".equals(type)) { // Respond to ping conn.send("{\"type\":\"pong\"}"); } } catch (JSONException e) { e.printStackTrace(); } }
@Override public void onError(WebSocket conn, Exception ex) { ex.printStackTrace(); }
public void broadcast(String message) { // Broadcast to all clients for (WebSocket client : clients) { client.send(message); } }
public static void main(String[] args) { ChatServer server = new ChatServer(8765); server.start(); }}import WebSocket from 'ws';
class WebSocketServer { private clients: Set<WebSocket> = new Set();
constructor(private port: number) {}
start(): void { const wss = new WebSocket.Server({ port: this.port });
wss.on('connection', (ws: WebSocket) => { this.registerClient(ws);
ws.on('message', (message: string) => { try { const data = JSON.parse(message); this.processMessage(ws, data); } catch (error) { console.error('Error parsing message:', error); } });
ws.on('close', () => { this.unregisterClient(ws); }); }); }
private registerClient(ws: WebSocket): void { this.clients.add(ws); console.log(`Client connected. Total: ${this.clients.size}`); }
private unregisterClient(ws: WebSocket): void { this.clients.delete(ws); console.log(`Client disconnected. Total: ${this.clients.size}`); }
private processMessage(ws: WebSocket, data: any): void { const messageType = data.type;
if (messageType === 'chat') { this.broadcast({ type: 'chat', user: data.user, message: data.message, timestamp: data.timestamp }); } else if (messageType === 'ping') { ws.send(JSON.stringify({ type: 'pong' })); } }
private broadcast(message: any): void { const messageJson = JSON.stringify(message); this.clients.forEach((client) => { if (client.readyState === WebSocket.OPEN) { client.send(messageJson); } }); }}
const server = new WebSocketServer(8765);server.start();#include <websocketpp/config/asio_no_tls.hpp>#include <websocketpp/server.hpp>#include <set>#include <iostream>
typedef websocketpp::server<websocketpp::config::asio> server;typedef server::message_ptr message_ptr;
class WebSocketServer {private: server m_server; std::set<websocketpp::connection_hdl, std::owner_less<websocketpp::connection_hdl>> m_clients;
void on_open(websocketpp::connection_hdl hdl) { m_clients.insert(hdl); std::cout << "Client connected. Total: " << m_clients.size() << std::endl; }
void on_close(websocketpp::connection_hdl hdl) { m_clients.erase(hdl); std::cout << "Client disconnected. Total: " << m_clients.size() << std::endl; }
void on_message(websocketpp::connection_hdl hdl, message_ptr msg) { // Process message std::string payload = msg->get_payload(); // Parse JSON and handle message broadcast(payload); }
void broadcast(const std::string& message) { for (auto it : m_clients) { m_server.send(it, message, websocketpp::frame::opcode::text); } }
public: WebSocketServer() { m_server.init_asio(); m_server.set_open_handler(bind(&WebSocketServer::on_open, this, ::_1)); m_server.set_close_handler(bind(&WebSocketServer::on_close, this, ::_1)); m_server.set_message_handler(bind(&WebSocketServer::on_message, this, ::_1, ::_2)); }
void run(uint16_t port) { m_server.listen(port); m_server.start_accept(); m_server.run(); }};
int main() { WebSocketServer server; server.run(8765); return 0;}using System;using System.Collections.Generic;using System.Net.WebSockets;using System.Threading;using System.Threading.Tasks;
public class WebSocketServer { private readonly HashSet<WebSocket> clients = new HashSet<WebSocket>();
public async Task HandleClient(WebSocket webSocket) { RegisterClient(webSocket);
try { var buffer = new byte[1024 * 4]; while (webSocket.State == WebSocketState.Open) { var result = await webSocket.ReceiveAsync( new ArraySegment<byte>(buffer), CancellationToken.None );
if (result.MessageType == WebSocketMessageType.Text) { var message = System.Text.Encoding.UTF8.GetString(buffer, 0, result.Count); await ProcessMessage(webSocket, message); } else if (result.MessageType == WebSocketMessageType.Close) { await webSocket.CloseAsync( WebSocketCloseStatus.NormalClosure, "Closed by client", CancellationToken.None ); } } } catch (Exception ex) { Console.WriteLine($"Error: {ex.Message}"); } finally { UnregisterClient(webSocket); } }
private void RegisterClient(WebSocket ws) { clients.Add(ws); Console.WriteLine($"Client connected. Total: {clients.Count}"); }
private void UnregisterClient(WebSocket ws) { clients.Remove(ws); Console.WriteLine($"Client disconnected. Total: {clients.Count}"); }
private async Task ProcessMessage(WebSocket ws, string message) { // Parse JSON and handle message // Broadcast to all clients await Broadcast(message); }
private async Task Broadcast(string message) { var buffer = System.Text.Encoding.UTF8.GetBytes(message); var tasks = new List<Task>();
foreach (var client in clients) { if (client.State == WebSocketState.Open) { tasks.Add(client.SendAsync( new ArraySegment<byte>(buffer), WebSocketMessageType.Text, true, CancellationToken.None )); } }
await Task.WhenAll(tasks); }}class WebSocketClient { constructor(url) { this.url = url; this.ws = null; this.reconnectInterval = 1000; this.maxReconnectAttempts = 5; this.reconnectAttempts = 0; }
connect() { this.ws = new WebSocket(this.url);
this.ws.onopen = () => { console.log('WebSocket connected'); this.reconnectAttempts = 0; this.onOpen(); };
this.ws.onmessage = (event) => { const data = JSON.parse(event.data); this.onMessage(data); };
this.ws.onerror = (error) => { console.error('WebSocket error:', error); this.onError(error); };
this.ws.onclose = () => { console.log('WebSocket closed'); this.onClose(); this.reconnect(); }; }
send(message) { if (this.ws && this.ws.readyState === WebSocket.OPEN) { this.ws.send(JSON.stringify(message)); } else { console.error('WebSocket not connected'); } }
reconnect() { if (this.reconnectAttempts < this.maxReconnectAttempts) { this.reconnectAttempts++; setTimeout(() => { console.log(`Reconnecting... (${this.reconnectAttempts}/${this.maxReconnectAttempts})`); this.connect(); }, this.reconnectInterval * this.reconnectAttempts); } }
onOpen() { // Override in subclass }
onMessage(data) { // Override in subclass }
onError(error) { // Override in subclass }
onClose() { // Override in subclass }
close() { if (this.ws) { this.ws.close(); } }}
// Usageconst client = new WebSocketClient('ws://localhost:8765');
client.onMessage = (data) => { if (data.type === 'chat') { console.log(`${data.user}: ${data.message}`); }};
client.connect();
// Send messageclient.send({ type: 'chat', user: 'John', message: 'Hello!', timestamp: Date.now()});import asyncioimport websocketsimport json
class WebSocketClient: def __init__(self, url): self.url = url self.ws = None self.reconnect_interval = 1 self.max_reconnect_attempts = 5 self.reconnect_attempts = 0
async def connect(self): try: self.ws = await websockets.connect(self.url) print('WebSocket connected') self.reconnect_attempts = 0 await self.on_open() await self.listen() except Exception as e: print(f'Connection error: {e}') await self.reconnect()
async def listen(self): try: async for message in self.ws: data = json.loads(message) await self.on_message(data) except websockets.exceptions.ConnectionClosed: print('WebSocket closed') await self.on_close() await self.reconnect()
async def send(self, message): if self.ws and self.ws.open: await self.ws.send(json.dumps(message)) else: print('WebSocket not connected')
async def reconnect(self): if self.reconnect_attempts < self.max_reconnect_attempts: self.reconnect_attempts += 1 await asyncio.sleep(self.reconnect_interval * self.reconnect_attempts) print(f'Reconnecting... ({self.reconnect_attempts}/{self.max_reconnect_attempts})') await self.connect()
async def on_open(self): pass
async def on_message(self, data): pass
async def on_close(self): pass
async def close(self): if self.ws: await self.ws.close()
# Usageasync def main(): client = WebSocketClient('ws://localhost:8765')
async def handle_message(data): if data.get('type') == 'chat': print(f"{data.get('user')}: {data.get('message')}")
client.on_message = handle_message await client.connect()
# Send message await client.send({ 'type': 'chat', 'user': 'John', 'message': 'Hello!', 'timestamp': int(time.time()) })
asyncio.run(main())import org.java_websocket.client.WebSocketClient;import org.java_websocket.handshake.ServerHandshake;import java.net.URI;
public class ChatClient extends WebSocketClient { private int reconnectAttempts = 0; private final int maxReconnectAttempts = 5;
public ChatClient(URI serverUri) { super(serverUri); }
@Override public void onOpen(ServerHandshake handshake) { System.out.println("WebSocket connected"); reconnectAttempts = 0; onOpen(); }
@Override public void onMessage(String message) { // Parse JSON and handle message System.out.println("Received: " + message); onMessage(message); }
@Override public void onClose(int code, String reason, boolean remote) { System.out.println("WebSocket closed"); onClose(); reconnect(); }
@Override public void onError(Exception ex) { ex.printStackTrace(); reconnect(); }
public void sendMessage(String type, String user, String message) { String json = String.format( "{\"type\":\"%s\",\"user\":\"%s\",\"message\":\"%s\",\"timestamp\":%d}", type, user, message, System.currentTimeMillis() ); send(json); }
private void reconnect() { if (reconnectAttempts < maxReconnectAttempts) { reconnectAttempts++; try { Thread.sleep(1000 * reconnectAttempts); System.out.println("Reconnecting... (" + reconnectAttempts + "/" + maxReconnectAttempts + ")"); reconnect(); } catch (InterruptedException e) { e.printStackTrace(); } } }
protected void onOpen() {} protected void onMessage(String message) {} protected void onClose() {}
public static void main(String[] args) { try { ChatClient client = new ChatClient(new URI("ws://localhost:8765")); client.connect();
// Send message client.sendMessage("chat", "John", "Hello!"); } catch (Exception e) { e.printStackTrace(); } }}import WebSocket from 'ws';
class WebSocketClient { private ws: WebSocket | null = null; private url: string; private reconnectInterval: number = 1000; private maxReconnectAttempts: number = 5; private reconnectAttempts: number = 0;
constructor(url: string) { this.url = url; }
connect(): void { this.ws = new WebSocket(this.url);
this.ws.on('open', () => { console.log('WebSocket connected'); this.reconnectAttempts = 0; this.onOpen(); });
this.ws.on('message', (data: WebSocket.Data) => { const message = JSON.parse(data.toString()); this.onMessage(message); });
this.ws.on('error', (error: Error) => { console.error('WebSocket error:', error); this.onError(error); });
this.ws.on('close', () => { console.log('WebSocket closed'); this.onClose(); this.reconnect(); }); }
send(message: any): void { if (this.ws && this.ws.readyState === WebSocket.OPEN) { this.ws.send(JSON.stringify(message)); } else { console.error('WebSocket not connected'); } }
private reconnect(): void { if (this.reconnectAttempts < this.maxReconnectAttempts) { this.reconnectAttempts++; setTimeout(() => { console.log(`Reconnecting... (${this.reconnectAttempts}/${this.maxReconnectAttempts})`); this.connect(); }, this.reconnectInterval * this.reconnectAttempts); } }
protected onOpen(): void {} protected onMessage(data: any): void { if (data.type === 'chat') { console.log(`${data.user}: ${data.message}`); } } protected onError(error: Error): void {} protected onClose(): void {}
close(): void { if (this.ws) { this.ws.close(); } }}
// Usageconst client = new WebSocketClient('ws://localhost:8765');client.connect();
client.send({ type: 'chat', user: 'John', message: 'Hello!', timestamp: Date.now()});#include <websocketpp/config/asio_client.hpp>#include <websocketpp/client.hpp>#include <iostream>
typedef websocketpp::client<websocketpp::config::asio_client> client;
class WebSocketClient {private: client m_client; websocketpp::connection_hdl m_hdl; std::string m_uri; int m_reconnect_attempts = 0; const int m_max_reconnect_attempts = 5;
void on_open(client* c, websocketpp::connection_hdl hdl) { std::cout << "WebSocket connected" << std::endl; m_reconnect_attempts = 0; m_hdl = hdl; }
void on_message(client* c, websocketpp::connection_hdl hdl, client::message_ptr msg) { std::cout << "Received: " << msg->get_payload() << std::endl; }
void on_close(client* c, websocketpp::connection_hdl hdl) { std::cout << "WebSocket closed" << std::endl; reconnect(); }
void reconnect() { if (m_reconnect_attempts < m_max_reconnect_attempts) { m_reconnect_attempts++; std::this_thread::sleep_for(std::chrono::seconds(m_reconnect_attempts)); std::cout << "Reconnecting... (" << m_reconnect_attempts << "/" << m_max_reconnect_attempts << ")" << std::endl; connect(); } }
public: WebSocketClient(const std::string& uri) : m_uri(uri) { m_client.init_asio(); m_client.set_open_handler(bind(&WebSocketClient::on_open, this, &m_client, ::_1)); m_client.set_message_handler(bind(&WebSocketClient::on_message, this, &m_client, ::_1, ::_2)); m_client.set_close_handler(bind(&WebSocketClient::on_close, this, &m_client, ::_1)); }
void connect() { websocketpp::lib::error_code ec; client::connection_ptr con = m_client.get_connection(m_uri, ec); if (ec) { std::cout << "Connection error: " << ec.message() << std::endl; return; } m_client.connect(con); m_client.run(); }
void send(const std::string& message) { m_client.send(m_hdl, message, websocketpp::frame::opcode::text); }};
int main() { WebSocketClient client("ws://localhost:8765"); client.connect(); return 0;}using System;using System.Net.WebSockets;using System.Text;using System.Threading;using System.Threading.Tasks;
public class WebSocketClient { private ClientWebSocket ws; private Uri serverUri; private int reconnectAttempts = 0; private const int maxReconnectAttempts = 5;
public WebSocketClient(string uri) { serverUri = new Uri(uri); }
public async Task ConnectAsync() { ws = new ClientWebSocket(); try { await ws.ConnectAsync(serverUri, CancellationToken.None); Console.WriteLine("WebSocket connected"); reconnectAttempts = 0; OnOpen(); await ReceiveLoop(); } catch (Exception ex) { Console.WriteLine($"Connection error: {ex.Message}"); await Reconnect(); } }
private async Task ReceiveLoop() { var buffer = new byte[1024 * 4]; while (ws.State == WebSocketState.Open) { var result = await ws.ReceiveAsync( new ArraySegment<byte>(buffer), CancellationToken.None );
if (result.MessageType == WebSocketMessageType.Text) { var message = Encoding.UTF8.GetString(buffer, 0, result.Count); OnMessage(message); } else if (result.MessageType == WebSocketMessageType.Close) { await ws.CloseAsync( WebSocketCloseStatus.NormalClosure, "Closed by server", CancellationToken.None ); } } OnClose(); await Reconnect(); }
public async Task SendAsync(string message) { if (ws != null && ws.State == WebSocketState.Open) { var buffer = Encoding.UTF8.GetBytes(message); await ws.SendAsync( new ArraySegment<byte>(buffer), WebSocketMessageType.Text, true, CancellationToken.None ); } }
private async Task Reconnect() { if (reconnectAttempts < maxReconnectAttempts) { reconnectAttempts++; await Task.Delay(1000 * reconnectAttempts); Console.WriteLine($"Reconnecting... ({reconnectAttempts}/{maxReconnectAttempts})"); await ConnectAsync(); } }
protected virtual void OnOpen() {} protected virtual void OnMessage(string message) { Console.WriteLine($"Received: {message}"); } protected virtual void OnClose() {}
public async Task CloseAsync() { if (ws != null) { await ws.CloseAsync( WebSocketCloseStatus.NormalClosure, "Closed by client", CancellationToken.None ); } }}
// Usagevar client = new WebSocketClient("ws://localhost:8765");await client.ConnectAsync();await client.SendAsync("{\"type\":\"chat\",\"user\":\"John\",\"message\":\"Hello!\"}");Client opens HTTP connection, server streams events:
GET /events HTTP/1.1Host: example.comAccept: text/event-streamCache-Control: no-cacheServer responds with stream:
HTTP/1.1 200 OKContent-Type: text/event-streamCache-Control: no-cacheConnection: keep-alive
event: messagedata: Hello World
event: updatedata: {"user": "John", "status": "online"}
event: messagedata: Goodbyefrom flask import Flask, Response, jsonifyimport jsonimport timeimport threading
app = Flask(__name__)
@app.route('/events')def stream_events(): """SSE endpoint""" def event_stream(): while True: # Send event data = { 'timestamp': time.time(), 'message': 'Server update' } yield f"event: update\ndata: {json.dumps(data)}\n\n"
time.sleep(1) # Send every second
return Response( event_stream(), mimetype='text/event-stream', headers={ 'Cache-Control': 'no-cache', 'Connection': 'keep-alive', 'X-Accel-Buffering': 'no' # Disable buffering in nginx } )
@app.route('/notify')def notify(): """Trigger notification""" # In real app, this would trigger event return jsonify({'status': 'notification sent'})
if __name__ == '__main__': app.run(threaded=True)import org.springframework.http.MediaType;import org.springframework.web.bind.annotation.*;import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;import java.io.IOException;import java.util.concurrent.Executors;import java.util.concurrent.ScheduledExecutorService;import java.util.concurrent.TimeUnit;
@RestControllerpublic class SSEServer { private final ScheduledExecutorService executor = Executors.newScheduledThreadPool(10);
@GetMapping(value = "/events", produces = MediaType.TEXT_EVENT_STREAM_VALUE) public SseEmitter streamEvents() { SseEmitter emitter = new SseEmitter(Long.MAX_VALUE);
executor.scheduleAtFixedRate(() -> { try { SseEmitter.SseEventBuilder event = SseEmitter.event() .name("update") .data(Map.of("timestamp", System.currentTimeMillis(), "message", "Server update"));
emitter.send(event); } catch (IOException e) { emitter.completeWithError(e); } }, 0, 1, TimeUnit.SECONDS);
emitter.onCompletion(() -> executor.shutdown()); emitter.onTimeout(() -> executor.shutdown());
return emitter; }}import express from 'express';
const app = express();
app.get('/events', (req, res) => { // Set SSE headers res.setHeader('Content-Type', 'text/event-stream'); res.setHeader('Cache-Control', 'no-cache'); res.setHeader('Connection', 'keep-alive'); res.setHeader('X-Accel-Buffering', 'no');
// Send events const interval = setInterval(() => { const data = { timestamp: Date.now(), message: 'Server update' };
res.write(`event: update\n`); res.write(`data: ${JSON.stringify(data)}\n\n`); }, 1000);
// Clean up on client disconnect req.on('close', () => { clearInterval(interval); res.end(); });});
app.listen(3000);#include <cpprest/http_listener.h>#include <cpprest/http_msg.h>#include <pplx/pplxtask.h>
class SSEServer {public: void handleEvents(web::http::http_request request) { web::http::http_response response; response.set_status_code(web::http::status_codes::OK); response.headers().add("Content-Type", "text/event-stream"); response.headers().add("Cache-Control", "no-cache"); response.headers().add("Connection", "keep-alive");
// Send events in a loop auto task = pplx::create_task([response]() { // In production, use proper async streaming // Simplified example return response; });
request.reply(response); }};using Microsoft.AspNetCore.Mvc;using System;using System.Threading;using System.Threading.Tasks;
[ApiController]public class SSEServerController : ControllerBase { [HttpGet("events")] public async Task GetEvents(CancellationToken cancellationToken) { Response.ContentType = "text/event-stream"; Response.Headers.Add("Cache-Control", "no-cache"); Response.Headers.Add("Connection", "keep-alive");
while (!cancellationToken.IsCancellationRequested) { var data = new { timestamp = DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(), message = "Server update" };
await Response.WriteAsync($"event: update\n"); await Response.WriteAsync($"data: {System.Text.Json.JsonSerializer.Serialize(data)}\n\n"); await Response.Body.FlushAsync();
await Task.Delay(1000, cancellationToken); } }}class SSEClient { constructor(url) { this.url = url; this.eventSource = null; }
connect() { this.eventSource = new EventSource(this.url);
this.eventSource.onopen = () => { console.log('SSE connected'); this.onOpen(); };
this.eventSource.onmessage = (event) => { const data = JSON.parse(event.data); this.onMessage(data); };
// Listen for specific event types this.eventSource.addEventListener('update', (event) => { const data = JSON.parse(event.data); this.onUpdate(data); });
this.eventSource.onerror = (error) => { console.error('SSE error:', error); this.onError(error); }; }
onOpen() { // Override in subclass }
onMessage(data) { // Override in subclass }
onUpdate(data) { // Override in subclass }
onError(error) { // Override in subclass }
close() { if (this.eventSource) { this.eventSource.close(); } }}
// Usageconst client = new SSEClient('/events');
client.onUpdate = (data) => { console.log('Update:', data);};
client.connect();class SSEClient { private url: string; private eventSource: EventSource | null = null;
constructor(url: string) { this.url = url; }
connect(): void { this.eventSource = new EventSource(this.url);
this.eventSource.onopen = () => { console.log('SSE connected'); this.onOpen(); };
this.eventSource.onmessage = (event: MessageEvent) => { const data = JSON.parse(event.data); this.onMessage(data); };
this.eventSource.addEventListener('update', (event: Event) => { const messageEvent = event as MessageEvent; const data = JSON.parse(messageEvent.data); this.onUpdate(data); });
this.eventSource.onerror = (error: Event) => { console.error('SSE error:', error); this.onError(error); }; }
protected onOpen(): void {} protected onMessage(data: any): void {} protected onUpdate(data: any): void { console.log('Update:', data); } protected onError(error: Event): void {}
close(): void { if (this.eventSource) { this.eventSource.close(); } }}
// Usageconst client = new SSEClient('/events');client.connect();import requestsimport json
class SSEClient: def __init__(self, url): self.url = url self.session = requests.Session()
def connect(self): response = self.session.get( self.url, stream=True, headers={'Accept': 'text/event-stream'} )
for line in response.iter_lines(): if line: line = line.decode('utf-8') if line.startswith('data: '): data = json.loads(line[6:]) self.on_message(data) elif line.startswith('event: '): event_type = line[7:] # Handle event type
def on_message(self, data): print(f"Received: {data}")
# Usageclient = SSEClient('http://localhost:3000/events')client.connect()import java.io.BufferedReader;import java.io.InputStreamReader;import java.net.HttpURLConnection;import java.net.URL;
public class SSEClient { private final String url;
public SSEClient(String url) { this.url = url; }
public void connect() { try { URL urlObj = new URL(url); HttpURLConnection connection = (HttpURLConnection) urlObj.openConnection(); connection.setRequestMethod("GET"); connection.setRequestProperty("Accept", "text/event-stream"); connection.setDoInput(true);
BufferedReader reader = new BufferedReader( new InputStreamReader(connection.getInputStream()) );
String line; while ((line = reader.readLine()) != null) { if (line.startsWith("data: ")) { String data = line.substring(6); onMessage(data); } } } catch (Exception e) { e.printStackTrace(); } }
protected void onMessage(String data) { System.out.println("Received: " + data); }}
// UsageSSEClient client = new SSEClient("http://localhost:3000/events");client.connect();#include <curl/curl.h>#include <string>#include <iostream>
class SSEClient {private: std::string url; CURL* curl;
static size_t WriteCallback(void* contents, size_t size, size_t nmemb, void* userp) { std::string* data = (std::string*)userp; data->append((char*)contents, size * nmemb); return size * nmemb; }
public: SSEClient(const std::string& url) : url(url) { curl = curl_easy_init(); }
void connect() { if (curl) { curl_easy_setopt(curl, CURLOPT_URL, url.c_str()); curl_easy_setopt(curl, CURLOPT_WRITEFUNCTION, WriteCallback);
struct curl_slist* headers = NULL; headers = curl_slist_append(headers, "Accept: text/event-stream"); curl_easy_setopt(curl, CURLOPT_HTTPHEADER, headers);
std::string response; curl_easy_setopt(curl, CURLOPT_WRITEDATA, &response);
CURLcode res = curl_easy_perform(curl); if (res == CURLE_OK) { onMessage(response); }
curl_slist_free_all(headers); } }
virtual void onMessage(const std::string& data) { std::cout << "Received: " << data << std::endl; }
~SSEClient() { if (curl) { curl_easy_cleanup(curl); } }};using System;using System.Net.Http;using System.Threading;using System.Threading.Tasks;
public class SSEClient { private readonly string url; private readonly HttpClient httpClient;
public SSEClient(string url) { this.url = url; this.httpClient = new HttpClient(); }
public async Task ConnectAsync(CancellationToken cancellationToken = default) { var request = new HttpRequestMessage(HttpMethod.Get, url); request.Headers.Add("Accept", "text/event-stream");
var response = await httpClient.SendAsync( request, HttpCompletionOption.ResponseHeadersRead, cancellationToken );
using (var stream = await response.Content.ReadAsStreamAsync()) { using (var reader = new System.IO.StreamReader(stream)) { string line; while ((line = await reader.ReadLineAsync()) != null) { if (line.StartsWith("data: ")) { var data = line.Substring(6); OnMessage(data); } } } } }
protected virtual void OnMessage(string data) { Console.WriteLine($"Received: {data}"); }}
// Usagevar client = new SSEClient("http://localhost:3000/events");await client.ConnectAsync();Client sends request, server holds it open:
from flask import Flask, jsonify, requestimport timeimport threading
app = Flask(__name__)pending_requests = []
@app.route('/poll')def poll(): """Long polling endpoint""" timeout = int(request.args.get('timeout', 30)) last_id = int(request.args.get('last_id', 0))
# Check for new data new_data = get_data_since(last_id)
if new_data: return jsonify(new_data)
# No data - hold request event = threading.Event() pending_requests.append({ 'event': event, 'last_id': last_id, 'timeout': timeout })
# Wait for data or timeout if event.wait(timeout): new_data = get_data_since(last_id) return jsonify(new_data) else: return jsonify({'status': 'timeout'}), 200
@app.route('/notify')def notify(): """Trigger notification""" # Wake up pending requests for req in pending_requests: req['event'].set()
pending_requests.clear() return jsonify({'status': 'notified'})import org.springframework.web.bind.annotation.*;import org.springframework.http.ResponseEntity;import java.util.concurrent.CompletableFuture;import java.util.concurrent.TimeUnit;
@RestControllerpublic class LongPollingServer {
@GetMapping("/poll") public CompletableFuture<ResponseEntity<Map<String, Object>>> poll( @RequestParam(defaultValue = "0") int lastId, @RequestParam(defaultValue = "30") int timeout) {
// Check for new data List<Data> newData = getDataSince(lastId);
if (!newData.isEmpty()) { return CompletableFuture.completedFuture( ResponseEntity.ok(Map.of("data", newData)) ); }
// No data - wait for new data or timeout return CompletableFuture.supplyAsync(() -> { try { // Wait for data or timeout Thread.sleep(timeout * 1000);
newData = getDataSince(lastId); if (!newData.isEmpty()) { return ResponseEntity.ok(Map.of("data", newData)); } else { return ResponseEntity.ok(Map.of("status", "timeout")); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); return ResponseEntity.ok(Map.of("status", "interrupted")); } }); }}import WebSocket from 'ws';import { EventEmitter } from 'events';
class ConnectionManager extends EventEmitter { private connections: Map<string, Set<WebSocket>> = new Map(); private heartbeatInterval: number = 30000;
addConnection(connectionId: string, ws: WebSocket): void { if (!this.connections.has(connectionId)) { this.connections.set(connectionId, new Set()); }
this.connections.get(connectionId)!.add(ws);
const interval = setInterval(() => { if (ws.readyState === WebSocket.OPEN) { ws.send(JSON.stringify({ type: 'ping' })); } else { clearInterval(interval); this.removeConnection(connectionId, ws); } }, this.heartbeatInterval);
ws.on('close', () => { clearInterval(interval); this.removeConnection(connectionId, ws); }); }
removeConnection(connectionId: string, ws: WebSocket): void { const conns = this.connections.get(connectionId); if (conns) { conns.delete(ws); if (conns.size === 0) { this.connections.delete(connectionId); } } }
async broadcast(message: any, connectionId?: string): Promise<void> { const targets = connectionId ? this.connections.get(connectionId) || new Set() : Array.from(this.connections.values()).flatMap(set => Array.from(set));
const messageJson = JSON.stringify(message); const disconnected: WebSocket[] = [];
for (const ws of targets) { if (ws.readyState === WebSocket.OPEN) { ws.send(messageJson); } else { disconnected.push(ws); } }
for (const ws of disconnected) { this.removeConnection('unknown', ws); } }}#include <map>#include <set>#include <websocketpp/server.hpp>#include <thread>#include <chrono>
class ConnectionManager {private: std::map<std::string, std::set<websocketpp::connection_hdl>> connections; int heartbeatInterval = 30;
public: void addConnection(const std::string& connectionId, websocketpp::connection_hdl hdl) { connections[connectionId].insert(hdl); }
void removeConnection(const std::string& connectionId, websocketpp::connection_hdl hdl) { auto it = connections.find(connectionId); if (it != connections.end()) { it->second.erase(hdl); if (it->second.empty()) { connections.erase(it); } } }
void broadcast(const std::string& message, const std::string& connectionId = "") { // Broadcast implementation }};using System;using System.Collections.Generic;using System.Linq;using System.Net.WebSockets;using System.Threading;using System.Threading.Tasks;
public class ConnectionManager { private readonly Dictionary<string, HashSet<WebSocket>> connections; private readonly int heartbeatInterval = 30000;
public ConnectionManager() { this.connections = new Dictionary<string, HashSet<WebSocket>>(); }
public void AddConnection(string connectionId, WebSocket webSocket) { lock (connections) { if (!connections.ContainsKey(connectionId)) { connections[connectionId] = new HashSet<WebSocket>(); } connections[connectionId].Add(webSocket); }
Task.Run(async () => { while (webSocket.State == WebSocketState.Open) { await Task.Delay(heartbeatInterval); if (webSocket.State == WebSocketState.Open) { var buffer = System.Text.Encoding.UTF8.GetBytes( System.Text.Json.JsonSerializer.Serialize(new { type = "ping" }) ); await webSocket.SendAsync( new ArraySegment<byte>(buffer), WebSocketMessageType.Text, true, CancellationToken.None ); } } RemoveConnection(connectionId, webSocket); }); }
public void RemoveConnection(string connectionId, WebSocket webSocket) { lock (connections) { if (connections.ContainsKey(connectionId)) { connections[connectionId].Remove(webSocket); if (connections[connectionId].Count == 0) { connections.Remove(connectionId); } } } }
public async Task BroadcastAsync(object message, string connectionId = null) { var messageJson = System.Text.Json.JsonSerializer.Serialize(message); var buffer = System.Text.Encoding.UTF8.GetBytes(messageJson);
var targets = connectionId != null ? connections.GetValueOrDefault(connectionId, new HashSet<WebSocket>()) : connections.Values.SelectMany(set => set).ToList();
var tasks = targets .Where(ws => ws.State == WebSocketState.Open) .Select(ws => ws.SendAsync( new ArraySegment<byte>(buffer), WebSocketMessageType.Text, true, CancellationToken.None ));
await Task.WhenAll(tasks); }}| Pattern | Use When | Don’t Use When |
|---|---|---|
| WebSocket | Bidirectional needed, low latency, high frequency | Simple one-way updates, HTTP-only environments |
| SSE | Server-to-client only, simple implementation, HTTP-based | Bidirectional needed, client-to-server messages |
| Long Polling | WebSocket/SSE not available, simple use case | High frequency, low latency needed |
Need bidirectional? Yes → WebSocket No → Need low latency? Yes → SSE No → Long Polling (fallback)Managing connections is critical:
from typing import Dict, Setimport asynciofrom datetime import datetime, timedelta
class ConnectionManager: """Manages WebSocket connections"""
def __init__(self): self.connections: Dict[str, Set] = {} self.heartbeat_interval = 30 # seconds self.connection_timeout = 60 # seconds
async def add_connection(self, connection_id: str, websocket): """Add new connection""" if connection_id not in self.connections: self.connections[connection_id] = set()
self.connections[connection_id].add(websocket)
# Start heartbeat asyncio.create_task(self.heartbeat(websocket))
async def remove_connection(self, connection_id: str, websocket): """Remove connection""" if connection_id in self.connections: self.connections[connection_id].discard(websocket) if not self.connections[connection_id]: del self.connections[connection_id]
async def heartbeat(self, websocket): """Send heartbeat to keep connection alive""" try: while True: await asyncio.sleep(self.heartbeat_interval) await websocket.send(json.dumps({'type': 'ping'})) except: await self.remove_connection('unknown', websocket)
async def broadcast(self, message, connection_id: str = None): """Broadcast message""" targets = self.connections.get(connection_id, set()) if connection_id else set().union(*self.connections.values())
disconnected = set() for ws in targets: try: await ws.send(json.dumps(message)) except: disconnected.add(ws)
# Clean up disconnected for ws in disconnected: await self.remove_connection('unknown', ws)import java.util.concurrent.ConcurrentHashMap;import java.util.Set;import java.util.concurrent.CopyOnWriteArraySet;import java.util.concurrent.ScheduledExecutorService;import java.util.concurrent.Executors;import java.util.concurrent.TimeUnit;
public class ConnectionManager { private final Map<String, Set<WebSocket>> connections = new ConcurrentHashMap<>(); private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(10); private final int heartbeatInterval = 30; // seconds
public void addConnection(String connectionId, WebSocket websocket) { connections.computeIfAbsent(connectionId, k -> new CopyOnWriteArraySet<>()) .add(websocket);
// Start heartbeat scheduler.scheduleAtFixedRate(() -> { try { websocket.send("{\"type\":\"ping\"}"); } catch (Exception e) { removeConnection(connectionId, websocket); } }, heartbeatInterval, heartbeatInterval, TimeUnit.SECONDS); }
public void removeConnection(String connectionId, WebSocket websocket) { Set<WebSocket> conns = connections.get(connectionId); if (conns != null) { conns.remove(websocket); if (conns.isEmpty()) { connections.remove(connectionId); } } }
public void broadcast(String message, String connectionId) { Set<WebSocket> targets = connectionId != null ? connections.getOrDefault(connectionId, Collections.emptySet()) : connections.values().stream() .flatMap(Set::stream) .collect(Collectors.toSet());
for (WebSocket ws : targets) { try { ws.send(message); } catch (Exception e) { removeConnection(connectionId, ws); } } }}🔄 WebSocket: Full-Duplex
WebSocket provides bidirectional, persistent connection. Best for chat, gaming, real-time collaboration.
📡 SSE: Simple Streaming
SSE is HTTP-based, server-to-client streaming. Simpler than WebSocket, automatic reconnection.
⏳ Long Polling: Fallback
Long polling holds requests open. Use when WebSocket/SSE not available. Less efficient but works everywhere.
🔌 Connection Management
Manage connections carefully: heartbeat, cleanup, reconnection. Critical for production systems.