A simple WebSocket server implementation in Go with support for pub/sub channels, presence channels, and private channels inspired from Laravel Echo
- WebSocket server with pub/sub capabilities
- Support for private channels with authorization
- Presence channels with join/leave events
- Client-to-client whisper events within channels
- Redis backend for scaling across multiple nodes
- Automatic connection management and cleanup
- Built-in statistics tracking
- Extensible message broker and storage interfaces
- Context-based graceful shutdown
- Comprehensive security features
- Configurable rate limiting
- Middleware support
- Customizable connection settings
-
Install the package using
go get:go get github.com/vortechron/go-ws -
Import the package in your Go code:
import "github.com/vortechron/go-ws"
-
Create a configuration:
config := ws.Config{ PingInterval: 30 * time.Second, WriteBufferSize: 1024, ReadBufferSize: 1024, RateLimit: ws.RateLimit{ Messages: 100, Interval: time.Minute, }, ShouldLogStats: true, }
-
Create a new WebSocket server instance with context for graceful shutdown:
ctx := context.Background() redisAddr := "localhost:6379" hub := ws.NewHub( ctx, ws.NewRedisBroker(redisAddr), ws.NewRedisStorage(redisAddr), config, ) go hub.Run()
-
Define your WebSocket options including middleware, auth handlers, and error handlers:
options := &ws.Options{ // Authorization handler AuthHandler: func(userID string, channelName string) bool { // Implement your authorization logic here return true }, // User identification GetUserID: func(r *http.Request) string { // Extract user ID from JWT token token := jwt.ParseFromRequest(r) return token.Claims.UserID }, // Error handling ErrorHandler: func(err error) { log.Printf("WebSocket error: %v", err) }, // Middleware chain Middleware: []ws.Middleware{ // ws.RateLimitMiddleware(), example }, // Connection events OnConnect: func(client *ws.Client) { log.Printf("Client connected: %s", client.ID) }, OnDisconnect: func(client *ws.Client) { log.Printf("Client disconnected: %s", client.ID) }, // Enable whisper functionality EnableWhispers: true, // Optional whisper middleware for filtering or transforming whispers WhisperMiddleware: func(whisper *ws.WhisperEvent) bool { // Process the whisper event here // Return false to block the whisper return true }, }
-
Start the WebSocket server with graceful shutdown:
// Create server with context server := &http.Server{ Addr: ":8080", } // Handle WebSocket connections http.HandleFunc("/ws", func(w http.ResponseWriter, r *http.Request) { if err := ws.ServeWS(hub, w, r, options); err != nil { log.Printf("Failed to serve WebSocket: %v", err) } }) // Start server go func() { if err := server.ListenAndServeTLS("cert.pem", "key.pem"); err != nil { log.Printf("Server error: %v", err) } }() // Graceful shutdown stop := make(chan os.Signal, 1) signal.Notify(stop, os.Interrupt) <-stop shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() if err := server.Shutdown(shutdownCtx); err != nil { log.Printf("Error during shutdown: %v", err) }
-
Connect and use from client side:
const socket = new WebSocket("wss://localhost:8080/ws"); // Handle connection socket.onopen = function() { // Subscribe to channels socket.send(JSON.stringify({ action: "subscribe", channel: "my-channel" })); // Subscribe to presence channel socket.send(JSON.stringify({ action: "subscribe", channel: "presence-room-1" })); }; // Handle incoming messages socket.onmessage = function(event) { const data = JSON.parse(event.data); // Handle different message types switch(data.type) { case "message": console.log("New message:", data.message); break; case "presence": console.log("Presence update:", data.users); break; case "whisper": console.log(`Whisper from ${data.from}, event: ${data.event}, channel: ${data.channel}`); console.log("Whisper data:", data.data); break; case "error": console.error("Error:", data.error); break; } }; // Handle errors socket.onerror = function(error) { console.error("WebSocket error:", error); }; // Handle disconnection socket.onclose = function() { console.log("Connection closed"); // Implement reconnection logic here };
Whisper events allow clients to send events directly to other clients who are subscribed to the same channel without going through your application server. This is useful for features like typing indicators, read receipts, or any temporary client state that doesn't need to be stored.
-
Send a whisper to a channel:
// Send a whisper to a channel socket.send(JSON.stringify({ action: "whisper", channel: "chat-room-1", // Channel to whisper on event: "typing", // Custom event name data: { // Custom data payload isTyping: true } }));
-
Listen for whispers on the client:
// Using the onmessage handler from above socket.onmessage = function(event) { const data = JSON.parse(event.data); if (data.type === "whisper") { // Handle specific whisper events switch(data.event) { case "typing": const isTyping = data.data.isTyping; console.log(`User ${data.from} is ${isTyping ? 'typing' : 'stopped typing'} in ${data.channel}`); updateTypingIndicator(data.from, isTyping, data.channel); break; case "read": console.log(`User ${data.from} read your message in ${data.channel}`); updateReadStatus(data.from, data.channel); break; default: console.log(`Received whisper from ${data.from} with event: ${data.event} in ${data.channel}`); console.log("Data:", data.data); } } };
-
Server-side configuration (Go):
// Enable whisper functionality in options options.EnableWhispers = true // Optional: Add whisper middleware for filtering or transforming whispers options.WhisperMiddleware = func(whisper *ws.WhisperEvent) bool { // You can modify the whisper event here // Return false to block the whisper // Example: Log all whispers log.Printf("Whisper from %s on channel %s: %s", whisper.FromID, whisper.ChannelName, whisper.Event) // Example: Block whispers with certain events if whisper.Event == "blocked-event" { return false } return true }
The package provides comprehensive error handling:
// Custom error handling
options.ErrorHandler = func(err error) {
switch e := err.(type) {
case *ws.AuthError:
// Handle authentication errors
case *ws.RateLimitError:
// Handle rate limiting errors
case *ws.ConnectionError:
// Handle connection errors
case *ws.WhisperError:
// Handle whisper-related errors
log.Printf("Whisper error: %s", e.Message)
default:
// Handle other errors
}
}When scaling across multiple nodes:
- Ensure Redis cluster is properly configured
- Configure appropriate connection limits per node
- Monitor memory usage and connection counts
- Use health checks for load balancing
// Configure connection limits
config.MaxConnections = 10000
config.MaxConnectionsPerIP = 100- Always use SSL/TLS in production
- Implement rate limiting
- Validate origin headers
- Use token-based authentication
- Sanitize all input messages
- Set appropriate timeouts
The package provides interfaces for extending functionality:
MessageBroker: Implement this interface to use a different message brokerStorageClient: Implement this interface to use a different storage backendMiddleware: Implement this interface to add custom middlewareAuthProvider: Implement this interface to add custom authentication
- Set up monitoring and alerting
- Configure appropriate timeouts
- Implement rate limiting
- Set up error tracking
- Configure logging
- Set up health checks
- Plan for scaling
- Implement reconnection strategy
- Set up backup and recovery procedures
Contributions are welcome! Please feel free to submit a Pull Request.
This project is licensed under the MIT License - see the LICENSE file for details.
If you're using the Fiber web framework, you can use the provided Fiber adapter:
- First, install the required dependencies:
go get github.com/gofiber/fiber/v2
go get github.com/gofiber/websocket/v2
- Create a new Fiber app and set up the WebSocket endpoint:
package main
import (
"context"
"time"
"github.com/gofiber/fiber/v2"
"github.com/vortechron/go-ws"
"github.com/vortechron/go-ws/adaptor"
)
func main() {
// Create WebSocket hub
ctx := context.Background()
redisAddr := "localhost:6379"
config := ws.Config{
PingInterval: 30 * time.Second,
WriteBufferSize: 1024,
ReadBufferSize: 1024,
}
hub, _ := ws.NewWebSocketServer(ctx, redisAddr, config)
go hub.Run()
// Create WebSocket options
options := &ws.Options{
// Authorization handler
AuthHandler: func(userID string, channelName string) bool {
// Implement your authorization logic
return true
},
// Get user ID from request
GetUserID: func(r *http.Request) string {
// Get user ID from request
return "user-123"
},
}
// Create Fiber app
app := fiber.New()
// Create WebSocket route
app.Use("/ws", adaptor.FiberHandler(hub, options))
// Start server
app.Listen(":8080")
}- Connect from the client side as shown in the earlier examples.