// client2server - Go server with WebSocket + Redpanda // Copyright (c) 2026 Luis Rosales - MIT License // // Build: go build -o client2server-server // Run: ./client2server-server // WebSocket: ws://localhost:3843 // HTTP API: http://localhost:3843 package main import ( "context" "encoding/json" "fmt" "log" "net/http" "os" "os/signal" "syscall" "time" "github.com/google/uuid" "github.com/redpanda-data/redpanda-sdk-go/redpanda" "github.com/redpanda-data/redpanda-sdk-go/schema" "nhooyr.io/websocket" ) // Config type Config struct { RedpandaBrokers []string WebsocketPort int APIPort int Token string } var cfg = Config{ RedpandaBrokers: []string{"localhost:9092"}, WebsocketPort: 3843, APIPort: 3844, // HTTP API on next port Token: "secret-token", } // Event from router type RouterEvent struct { ID string `json:"id"` RouterID string `json:"router_id"` Hostname string `json:"hostname"` EventType string `json:"event_type"` Timestamp time.Time `json:"timestamp"` Payload map[string]interface{} `json:"payload"` ReceivedAt time.Time `json:"received_at"` Connection string `json:"connection"` // "websocket" or "http" } // Command to router type RouterCommand struct { ID string `json:"id"` RouterID string `json:"router_id"` Command string `json:"command"` Args map[string]string `json:"args,omitempty"` SentAt time.Time `json:"sent_at"` } // Router state type Router struct { ID string LastSeen time.Time Conn *websocket.Conn Connected bool } var routers = make(map[string]*Router) // Redpanda var rp *redpanda.Client func initRedpanda() error { cfg.RedpandaBrokers = getEnvComma("REDPANDA_BROKERS", "localhost:9092") var err error rp, err = redpanda.NewClient(&redpanda.ClientConfig{ Brokers: cfg.RedpandaBrokers, }) if err != nil { return fmt.Errorf("redpanda: %v", err) } // Create topics if needed topics := []string{"router-events", "router-commands"} for _, topic := range topics { err := rp.CreateTopic(topic, 1, 3) if err != nil && !schema.ErrTopicExists.Exists(err) { log.Printf("Topic %s: %v", topic, err) } } return nil } func getEnvComma(key, def string) []string { val := os.Getenv(key) if val == "" { return []string{def} } return []string{val} } // Publish event to Redpanda func publishEvent(event RouterEvent) error { data, err := json.Marshal(event) if err != nil { return err } return rp.Produce("router-events", []byte(event.ID), data) } // Publish command to Redpanda func publishCommand(cmd RouterCommand) error { data, err := json.Marshal(cmd) if err != nil { return err } return rp.Produce("router-commands", []byte(cmd.ID), data) } // WebSocket handler func handleWebSocket(w http.ResponseWriter, r *http.Request) { // Auth token := r.URL.Query().Get("token") if token != cfg.Token { http.Error(w, "Unauthorized", http.StatusUnauthorized) return } conn, err := websocket.Accept(w, r, &websocket.AcceptOptions{ CompressionMode: websocket.CompressionContextTakeover, }) if err != nil { log.Printf("WS accept: %v", err) return } defer conn.Close(websocket.StatusNormalClosure, "") ctx := context.Background() // Read router registration var regMsg RouterEvent err = conn.Read(ctx, ®Msg) if err != nil { log.Printf("WS read reg: %v", err) return } routerID := regMsg.RouterID if routerID == "" { routerID = r.RemoteAddr } // Store router routers[routerID] = &Router{ ID: routerID, LastSeen: time.Now(), Conn: conn, Connected: true, } log.Printf("Router connected: %s", routerID) // Update router last seen routers[routerID].LastSeen = time.Now() // Message loop for { var event RouterEvent err := conn.Read(ctx, &event) if err != nil { break } event.ID = uuid.New().String() event.ReceivedAt = time.Now() event.Connection = "websocket" // Publish to Redpanda if err := publishEvent(event); err != nil { log.Printf("Publish error: %v", err) } log.Printf("[%s] %s", routerID, event.EventType) // Send ack conn.Write(ctx, []byte(`{"ack":true}`)) // Update last seen routers[routerID].LastSeen = time.Now() } // Cleanup if routers[routerID] != nil { routers[routerID].Connected = false } log.Printf("Router disconnected: %s", routerID) } // HTTP Event webhook (fallback) func handleHTTPEvent(w http.ResponseWriter, r *http.Request) { if r.Method != "POST" { http.Error(w, "Method not allowed", http.StatusMethodNotAllowed) return } token := r.Header.Get("Authorization") if token != "Bearer "+cfg.Token && token != cfg.Token { http.Error(w, "Unauthorized", http.StatusUnauthorized) return } var event RouterEvent if err := json.NewDecoder(r.Body).Decode(&event); err != nil { http.Error(w, "Invalid JSON", http.StatusBadRequest) return } event.ID = uuid.New().String() event.ReceivedAt = time.Now() event.Connection = "http" // Publish to Redpanda if err := publishEvent(event); err != nil { log.Printf("Publish error: %v", err) w.WriteHeader(http.StatusInternalServerError) return } log.Printf("[%s] %s (HTTP)", event.RouterID, event.EventType) json.NewEncoder(w).Encode(map[string]string{"event_id": event.ID}) } // Get routers list func handleRouters(w http.ResponseWriter, r *http.Request) { json.NewEncoder(w).Encode(map[string]interface{}{ "routers": func() []map[string]interface{} { list := []map[string]interface{}{} for id, router := range routers { list = append(list, map[string]interface{}{ "id": id, "last_seen": router.LastSeen, "online": time.Since(router.LastSeen) < 60*time.Second, }) } return list }(), }) } // Get events func handleEvents(w http.ResponseWriter, r *http.Request) { // Simple consumer - in production use proper offset management consumer, err := rp.NewConsumer("router-events", "client2server-group") if err != nil { http.Error(w, err.Error(), http.StatusInternalServerError) return } defer consumer.Close() limit := 100 if l := r.URL.Query().Get("limit"); l != "" { fmt.Sscanf(l, "%d", &limit) } events := []RouterEvent{} ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() for i := 0; i < limit; i++ { msg, err := consumer.Consume(ctx, 5*time.Second) if err != nil { break } var event RouterEvent if json.Unmarshal(msg.Value, &event) == nil { events = append(events, event) } } json.NewEncoder(w).Encode(map[string]interface{}{ "events": events, "total": len(events), }) } // Send command to router func handleCommand(w http.ResponseWriter, r *http.Request) { if r.Method != "POST" { http.Error(w, "Method not allowed", http.StatusMethodNotAllowed) return } var cmd RouterCommand if err := json.NewDecoder(r.Body).Decode(&cmd); err != nil { http.Error(w, "Invalid JSON", http.StatusBadRequest) return } routerID := r.URL.Query().Get("router_id") if routerID == "" { routerID = cmd.RouterID } if routerID == "" { http.Error(w, "router_id required", http.StatusBadRequest) return } // Publish command cmd.ID = uuid.New().String() cmd.RouterID = routerID cmd.SentAt = time.Now() if err := publishCommand(cmd); err != nil { log.Printf("Publish command: %v", err) http.Error(w, err.Error(), http.StatusInternalServerError) return } log.Printf("Command to %s: %s", routerID, cmd.Command) json.NewEncoder(w).Encode(map[string]string{ "command_id": cmd.ID, "status": "sent", }) } // Health check func handleHealth(w http.ResponseWriter, r *http.Request) { json.NewEncoder(w).Encode(map[string]interface{}{ "status": "ok", "routers": len(routers), "timestamp": time.Now(), }) } // Main func main() { // Config from env cfg.Token = getEnvStr("TOKEN", cfg.Token) cfg.WebsocketPort = getEnvInt("PORT", cfg.WebsocketPort) log.SetFlags(log.LstdFlags | log.Lshortfile) log.Printf("=== client2server Go Server ===") log.Printf("WebSocket: ws://localhost:%d", cfg.WebsocketPort) log.Printf("HTTP API: http://localhost:%d", cfg.APIPort) // Init Redpanda if err := initRedpanda(); err != nil { log.Printf("Redpanda init failed (running without): %v", err) } else { log.Printf("Redpanda: %v", cfg.RedpandaBrokers) } // HTTP handlers http.HandleFunc("/", handleHealth) http.HandleFunc("/health", handleHealth) http.HandleFunc("/api/events", handleHTTPEvent) http.HandleFunc("/api/events/", handleHTTPEvent) http.HandleFunc("/api/routers", handleRouters) http.HandleFunc("/api/command", handleCommand) // WebSocket (HAProxy protocol) http.HandleFunc("/ws", handleWebSocket) // Graceful shutdown go func() { sigCh := make(chan os.Signal, 1) signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM) <-sigCh log.Println("Shutting down...") for _, router := range routers { router.Conn.Close(websocket.StatusNormalClosure, "server shutdown") } os.Exit(0) }() // Start (using HTTP for both WS and API - need haproxy or prefix) log.Printf("Server ready on port %d", cfg.WebiscrollPoint) log.Fatal(http.ListenAndServe(fmt.Sprintf(":%d", cfg.WebsocketPort), nil)) } func getEnvStr(key, def string) string { if val := os.Getenv(key); val != "" { return val } return def } func getEnvInt(key string, def int) int { if val := os.Getenv(key); val != "" { var v int fmt.Sscanf(val, "%d", &v) return v } return def }