// client2server - Go server with WebSocket + Redpanda + Dashboard API // Copyright (c) 2026 Luis Rosales - MIT License // // Stack: // - WebSocket: github.com/coder/websocket // - Kafka: github.com/twmb/franz-go (talks to Redpanda) // - Auth: JWT (HS256) with role-based access // - Storage: SQLite (modernc.org/sqlite, pure Go, no cgo) // - Live feed: Server-Sent Events // // Endpoints: // POST /api/auth/login - username/password -> JWT // GET /api/auth/me - current user // GET /health - liveness // GET /api/routers - known routers (legacy shared token OK) // GET /api/events?limit=&router_id= - event history (auth) // POST /api/events - ingest event (router or hotplug) // POST /api/command - send command to router (auth) // GET /api/commands?router_id= - command history // GET /api/metrics?since=1h - time-series metrics // GET /api/alerts?unack=1 - alerts // POST /api/alerts/:id/ack - acknowledge alert // GET /api/events/stream - SSE live feed // GET /ws - WebSocket from routers (legacy token) // // Env: REDPANDA_BROKERS, TOKEN (legacy shared), JWT_SECRET, DB_PATH, PORT package main import ( "context" "encoding/json" "errors" "fmt" "log" "net/http" "os" "os/signal" "strconv" "strings" "sync" "syscall" "time" "github.com/coder/websocket" "github.com/google/uuid" "github.com/twmb/franz-go/pkg/kgo" ) // ---------------------------------------------------------------------------- // Config // ---------------------------------------------------------------------------- type Config struct { RedpandaBrokers []string Port int Token string // legacy shared token for routers/hotplug DBPath string } var cfg = Config{ RedpandaBrokers: []string{"localhost:9092"}, Port: 3843, Token: "secret-token", DBPath: "data/client2server.db", } func loadConfigFromEnv() { cfg.RedpandaBrokers = strings.Split(getEnvStr("REDPANDA_BROKERS", "localhost:9092"), ",") cfg.Port = getEnvInt("PORT", cfg.Port) cfg.Token = getEnvStr("TOKEN", cfg.Token) cfg.DBPath = getEnvStr("DB_PATH", cfg.DBPath) } func getEnvStr(key, def string) string { if v := os.Getenv(key); v != "" { return v } return def } func getEnvInt(key string, def int) int { if v := os.Getenv(key); v != "" { var n int if _, err := fmt.Sscanf(v, "%d", &n); err == nil { return n } } return def } // ---------------------------------------------------------------------------- // Domain types // ---------------------------------------------------------------------------- type RouterEvent struct { ID string `json:"id"` RouterID string `json:"router_id"` Hostname string `json:"hostname,omitempty"` 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"` } 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"` } type CommandResult struct { Success bool `json:"success"` Output string `json:"output"` Error string `json:"error"` } // ---------------------------------------------------------------------------- // Router registry // ---------------------------------------------------------------------------- type Router struct { ID string LastSeen time.Time Conn *websocket.Conn Connected bool writeMu sync.Mutex } var ( routersMu sync.RWMutex routers = make(map[string]*Router) routerQueuesMu sync.Mutex routerQueues = make(map[string][]RouterCommand) pendingMu sync.Mutex pendingCmds = make(map[string]chan CommandResult) executedMu sync.Mutex executedCmds = make(map[string]time.Time) ) const ( commandTimeout = 30 * time.Second idempotencyTTL = 5 * time.Minute wsWriteTimeout = 10 * time.Second publishTimeout = 3 * time.Second offlineThreshold = 60 * time.Second ) // ---------------------------------------------------------------------------- // Redpanda (Kafka) producer // ---------------------------------------------------------------------------- var kcl *kgo.Client func initRedpanda(ctx context.Context) error { cl, err := kgo.NewClient( kgo.SeedBrokers(cfg.RedpandaBrokers...), kgo.ClientID("client2server"), kgo.ProducerLinger(5*time.Millisecond), kgo.ProducerBatchCompression(kgo.SnappyCompression()), ) if err != nil { return fmt.Errorf("kafka client: %w", err) } kcl = cl return nil } func publish(ctx context.Context, topic string, key string, value any) error { data, err := json.Marshal(value) if err != nil { return fmt.Errorf("marshal: %w", err) } rec := &kgo.Record{Topic: topic, Key: []byte(key), Value: data} // Fire-and-forget: don't block the HTTP path. franz-go queues records // internally and flushes them in the background. If the broker is down, // the record will be retried until it expires (RecordDeliveryTimeout). // We still get error visibility via a callback for logging. pctx, cancel := context.WithTimeout(ctx, 100*time.Millisecond) defer cancel() res := kcl.ProduceSync(pctx, rec) if err := res.FirstErr(); err != nil { log.Printf("publish %s/%s: %v", topic, key, err) return err } return nil } // ---------------------------------------------------------------------------- // Event ingestion (called by both WS and HTTP) // ---------------------------------------------------------------------------- func ingestEvent(ev RouterEvent) { if ev.ID == "" { ev.ID = uuid.New().String() } if ev.ReceivedAt.IsZero() { ev.ReceivedAt = time.Now() } if ev.Timestamp.IsZero() { ev.Timestamp = ev.ReceivedAt } // Metrics globalMetrics.recordEvent(ev.EventType) // Persist (best-effort) if err := saveEvent(ev); err != nil { log.Printf("save event: %v", err) } // Publish to Redpanda if err := publish(context.Background(), "router-events", ev.RouterID, ev); err != nil { log.Printf("publish event: %v", err) } // Broadcast to SSE clients hub.broadcast("event", ev) } // ---------------------------------------------------------------------------- // WebSocket handler (routers connect here) // ---------------------------------------------------------------------------- func handleWebSocket(w http.ResponseWriter, r *http.Request) { // Auth: legacy TOKEN only (routers use shared secret, not JWT) token := r.URL.Query().Get("token") if token == "" { if h := r.Header.Get("Authorization"); strings.HasPrefix(h, "Bearer ") { token = strings.TrimPrefix(h, "Bearer ") } } if token != cfg.Token { http.Error(w, "Unauthorized", http.StatusUnauthorized) return } conn, err := websocket.Accept(w, r, &websocket.AcceptOptions{ CompressionMode: websocket.CompressionDisabled, }) if err != nil { log.Printf("ws accept: %v", err) return } defer conn.Close(websocket.StatusNormalClosure, "bye") ctx, cancel := context.WithCancel(r.Context()) defer cancel() // First message is registration raw, err := readRouterMessage(ctx, conn) if err != nil { log.Printf("ws read reg: %v", err) return } var reg RouterEvent if err := json.Unmarshal(raw, ®); err != nil { log.Printf("ws reg parse: %v", err) return } routerID := reg.RouterID if routerID == "" { routerID = r.RemoteAddr } routersMu.Lock() routers[routerID] = &Router{ ID: routerID, LastSeen: time.Now(), Conn: conn, Connected: true, } routersMu.Unlock() // Mark as back online (clears any offline alert) acknowledgeRouterAlerts(routerID) log.Printf("router connected: %s (from %s)", routerID, r.RemoteAddr) flushQueuedCommands(ctx, routerID) for { raw, err := readRouterMessage(ctx, conn) if err != nil { break } var ev RouterEvent if err := json.Unmarshal(raw, &ev); err != nil { log.Printf("[%s] bad event json: %v", routerID, err) continue } // Inherit router_id if missing if ev.RouterID == "" { ev.RouterID = routerID } ev.Connection = "websocket" // Command result handling if ev.EventType == "command_result" { if cid, _ := ev.Payload["command_id"].(string); cid != "" { deliverCommandResult(cid, ev.Payload) executedMu.Lock() executedCmds[cid] = time.Now() executedMu.Unlock() } } ingestEvent(ev) log.Printf("[%s] %s", routerID, ev.EventType) routersMu.Lock() if r, ok := routers[routerID]; ok { r.LastSeen = time.Now() } routersMu.Unlock() } routersMu.Lock() if r, ok := routers[routerID]; ok { r.Connected = false } routersMu.Unlock() log.Printf("router disconnected: %s", routerID) } func readRouterMessage(ctx context.Context, conn *websocket.Conn) ([]byte, error) { _, data, err := conn.Read(ctx) return data, err } func deliverCommandResult(cmdID string, payload map[string]interface{}) { pendingMu.Lock() ch, ok := pendingCmds[cmdID] if ok { delete(pendingCmds, cmdID) } pendingMu.Unlock() if !ok { return } res := CommandResult{Success: payload["success"] == true} if s, ok := payload["output"].(string); ok { res.Output = s } if s, ok := payload["error"].(string); ok { res.Error = s } // Persist result to SQLite _ = updateCommandResult(cmdID, "completed", &res) globalMetrics.recordCommand(res.Success) select { case ch <- res: default: } } func flushQueuedCommands(ctx context.Context, routerID string) { routerQueuesMu.Lock() queue := routerQueues[routerID] delete(routerQueues, routerID) routerQueuesMu.Unlock() if len(queue) == 0 { return } log.Printf("flushing %d queued commands to %s", len(queue), routerID) for _, cmd := range queue { cmd.SentAt = time.Now() pendingMu.Lock() pendingCmds[cmd.ID] = make(chan CommandResult, 1) pendingMu.Unlock() _ = updateCommandResult(cmd.ID, "delivered", nil) if err := publish(ctx, "router-commands", routerID, cmd); err != nil { log.Printf("queue flush publish: %v", err) } } } func acknowledgeRouterAlerts(routerID string) { // Best-effort: mark all unacked offline alerts for this router as resolved if db == nil { return } _, _ = db.Exec("UPDATE alerts SET acknowledged_at = CURRENT_TIMESTAMP WHERE router_id = ? AND kind = 'router_offline' AND acknowledged_at IS NULL", routerID) } // ---------------------------------------------------------------------------- // HTTP API // ---------------------------------------------------------------------------- func handleHealth(w http.ResponseWriter, r *http.Request) { routersMu.RLock() onlineCount := 0 total := len(routers) for _, rt := range routers { if time.Since(rt.LastSeen) < offlineThreshold { onlineCount++ } } routersMu.RUnlock() w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(map[string]interface{}{ "status": "ok", "routers": total, "routers_online": onlineCount, "redpanda": cfg.RedpandaBrokers, "version": "2.1.0", }) } func handleRouters(w http.ResponseWriter, r *http.Request) { type entry struct { ID string `json:"id"` LastSeen time.Time `json:"last_seen"` Online bool `json:"online"` Queued int `json:"queued_commands"` } routersMu.RLock() list := make([]entry, 0, len(routers)) for id, rt := range routers { list = append(list, entry{ ID: id, LastSeen: rt.LastSeen, Online: time.Since(rt.LastSeen) < offlineThreshold, Queued: len(routerQueues[id]), }) } routersMu.RUnlock() w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(map[string]interface{}{"routers": list}) } func handleHTTPEvent(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { http.Error(w, "method not allowed", http.StatusMethodNotAllowed) return } // Accept either JWT or legacy shared TOKEN if !authorisedRouter(r) { http.Error(w, "unauthorized", http.StatusUnauthorized) return } var ev RouterEvent if err := json.NewDecoder(r.Body).Decode(&ev); err != nil { http.Error(w, "invalid json", http.StatusBadRequest) return } ev.Connection = "http" ingestEvent(ev) w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(map[string]string{ "event_id": ev.ID, "router_id": ev.RouterID, "status": "accepted", }) _ = ev } func handleListEvents(w http.ResponseWriter, r *http.Request) { limit, _ := strconv.Atoi(r.URL.Query().Get("limit")) if limit <= 0 || limit > 1000 { limit = 100 } routerID := r.URL.Query().Get("router_id") eventType := r.URL.Query().Get("event_type") events, err := listEvents(limit, routerID, eventType) if err != nil { http.Error(w, "db error: "+err.Error(), http.StatusInternalServerError) return } w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(map[string]interface{}{"events": events}) } func handleCommand(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { http.Error(w, "method not allowed", http.StatusMethodNotAllowed) return } // JWT or legacy TOKEN if !authorisedAny(r) { http.Error(w, "unauthorized", http.StatusUnauthorized) 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 } // Idempotency xecID := cmd.ID if xecID == "" { xecID = uuid.New().String() } executedMu.Lock() if last, exists := executedCmds[xecID]; exists && time.Since(last) < idempotencyTTL { executedMu.Unlock() w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(map[string]interface{}{ "idempotent_reject": true, "existing_command_id": xecID, "message": "command already executed within TTL", }) return } delete(executedCmds, xecID) executedMu.Unlock() cmd.ID = xecID cmd.RouterID = routerID cmd.SentAt = time.Now() // Identify the issuer (JWT username) if present issuedBy := "token" if claims := claimsFromHeader(r); claims != nil { issuedBy = claims.Username } // Persist _ = saveCommand(cmd, issuedBy, "pending") routersMu.RLock() router, online := routers[routerID] routersMu.RUnlock() if !online || router == nil || !router.Connected { routerQueuesMu.Lock() routerQueues[routerID] = append(routerQueues[routerID], cmd) depth := len(routerQueues[routerID]) routerQueuesMu.Unlock() log.Printf("queued cmd %s for %s (depth=%d)", cmd.Command, routerID, depth) _ = updateCommandResult(cmd.ID, "queued", nil) w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(map[string]interface{}{ "command_id": cmd.ID, "status": "queued", "queued_for": routerID, "queue_depth": depth, }) return } resultCh := make(chan CommandResult, 1) pendingMu.Lock() pendingCmds[cmd.ID] = resultCh pendingMu.Unlock() if err := publish(r.Context(), "router-commands", routerID, cmd); err != nil { pendingMu.Lock() delete(pendingCmds, cmd.ID) pendingMu.Unlock() _ = updateCommandResult(cmd.ID, "publish_failed", nil) http.Error(w, "publish failed: "+err.Error(), http.StatusBadGateway) return } _ = updateCommandResult(cmd.ID, "delivered", nil) log.Printf("cmd %s -> %s (waiting)", cmd.Command, routerID) w.Header().Set("Content-Type", "application/json") select { case res := <-resultCh: _ = json.NewEncoder(w).Encode(map[string]interface{}{ "command_id": cmd.ID, "status": "completed", "success": res.Success, "output": res.Output, "error": res.Error, }) case <-time.After(commandTimeout): pendingMu.Lock() delete(pendingCmds, cmd.ID) pendingMu.Unlock() _ = updateCommandResult(cmd.ID, "timeout", nil) _ = json.NewEncoder(w).Encode(map[string]string{ "command_id": cmd.ID, "status": "timeout", "error": "router did not respond", }) } } func handleListCommands(w http.ResponseWriter, r *http.Request) { limit, _ := strconv.Atoi(r.URL.Query().Get("limit")) if limit <= 0 || limit > 500 { limit = 50 } routerID := r.URL.Query().Get("router_id") cmds, err := listCommands(limit, routerID) if err != nil { http.Error(w, "db error: "+err.Error(), http.StatusInternalServerError) return } w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(map[string]interface{}{"commands": cmds}) } func handleMetrics(w http.ResponseWriter, r *http.Request) { sinceStr := r.URL.Query().Get("since") since := 1 * time.Hour if sinceStr != "" { if d, err := time.ParseDuration(sinceStr); err == nil { since = d } } w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(map[string]interface{}{ "summary": globalMetrics.summary(), "buckets": globalMetrics.snapshot(since), "since": since.String(), }) } func handleAlerts(w http.ResponseWriter, r *http.Request) { if r.Method == http.MethodPost { // Acknowledge: POST /api/alerts/{id}/ack // Path is set up by the mux (see main) http.Error(w, "use POST /api/alerts/{id}/ack", http.StatusMethodNotAllowed) return } unack := r.URL.Query().Get("unack") == "1" limit, _ := strconv.Atoi(r.URL.Query().Get("limit")) if limit <= 0 || limit > 500 { limit = 100 } alerts, err := listAlerts(unack, limit) if err != nil { http.Error(w, "db error: "+err.Error(), http.StatusInternalServerError) return } w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(map[string]interface{}{"alerts": alerts}) } func handleAckAlert(w http.ResponseWriter, r *http.Request) { if r.Method != http.MethodPost { http.Error(w, "method not allowed", http.StatusMethodNotAllowed) return } // Extract id from /api/alerts/{id}/ack parts := strings.Split(strings.Trim(r.URL.Path, "/"), "/") if len(parts) < 3 { http.Error(w, "bad path", http.StatusBadRequest) return } id, err := strconv.ParseInt(parts[2], 10, 64) if err != nil { http.Error(w, "bad id", http.StatusBadRequest) return } if err := acknowledgeAlert(id); err != nil { http.Error(w, err.Error(), http.StatusInternalServerError) return } w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(map[string]string{"status": "acknowledged"}) } func handleMe(w http.ResponseWriter, r *http.Request) { claims := claimsFromHeader(r) if claims == nil { http.Error(w, "unauthorized", http.StatusUnauthorized) return } w.Header().Set("Content-Type", "application/json") _ = json.NewEncoder(w).Encode(map[string]interface{}{ "user_id": claims.UserID, "username": claims.Username, "role": claims.Role, "expires": claims.ExpiresAt, }) } // ---------------------------------------------------------------------------- // Authorisation helpers // ---------------------------------------------------------------------------- func claimsFromHeader(r *http.Request) *Claims { auth := r.Header.Get("Authorization") token := trimBearer(auth) if token == "" || token == cfg.Token { return nil } claims, err := ParseJWT(token) if err != nil { return nil } return claims } func authorisedAny(r *http.Request) bool { auth := r.Header.Get("Authorization") token := trimBearer(auth) if token == "" { return false } if token == cfg.Token { return true } _, err := ParseJWT(token) return err == nil } func authorisedRouter(r *http.Request) bool { auth := r.Header.Get("Authorization") token := trimBearer(auth) if token == "" { return false } // Routers use the legacy shared TOKEN. JWT users are valid too (in case // someone scripts an event submission from the dashboard). if token == cfg.Token { return true } _, err := ParseJWT(token) return err == nil } // ---------------------------------------------------------------------------- // Background jobs // ---------------------------------------------------------------------------- func startJanitor(ctx context.Context) { go func() { t := time.NewTicker(time.Minute) defer t.Stop() for { select { case <-ctx.Done(): return case <-t.C: now := time.Now() executedMu.Lock() for id, ts := range executedCmds { if now.Sub(ts) > idempotencyTTL { delete(executedCmds, id) } } executedMu.Unlock() } } }() } func startOfflineWatcher(ctx context.Context) { go func() { t := time.NewTicker(30 * time.Second) defer t.Stop() known := map[string]bool{} for { select { case <-ctx.Done(): return case <-t.C: routersMu.RLock() current := map[string]bool{} for id, rt := range routers { online := time.Since(rt.LastSeen) < offlineThreshold current[id] = online if !online && !known[id] { createAlert(id, "router_offline", fmt.Sprintf("Router %s has been offline for %s", id, time.Since(rt.LastSeen).Round(time.Second))) } } routersMu.RUnlock() known = current } } }() } // ---------------------------------------------------------------------------- // main // ---------------------------------------------------------------------------- func main() { loadConfigFromEnv() log.SetFlags(log.LstdFlags | log.Lshortfile) log.Printf("=== client2server v2.1 ===") log.Printf("port=%d brokers=%v db=%s", cfg.Port, cfg.RedpandaBrokers, cfg.DBPath) // Ensure DB directory exists if err := os.MkdirAll(strings.TrimSuffix(cfg.DBPath, "/"+pathBase(cfg.DBPath)), 0755); err != nil { log.Printf("mkdir db: %v", err) } if err := initStore(cfg.DBPath); err != nil { log.Fatalf("init store: %v", err) } ctx, cancel := context.WithCancel(context.Background()) defer cancel() if err := initRedpanda(ctx); err != nil { log.Printf("redpanda init failed (continuing): %v", err) } else { defer kcl.Close() // Try to ensure topics exist (best effort) if err := ensureTopics(ctx); err != nil { log.Printf("ensure topics: %v", err) } // Start the consumer that persists to SQLite go startConsumer(ctx) } startJanitor(ctx) startOfflineWatcher(ctx) mux := http.NewServeMux() // Public mux.HandleFunc("/health", handleHealth) mux.HandleFunc("/", handleHealth) // Auth mux.HandleFunc("/api/auth/login", handleLogin) // Router-facing (legacy token OR JWT) mux.HandleFunc("/api/events", handleHTTPEvent) mux.HandleFunc("/api/events/", handleHTTPEvent) mux.HandleFunc("/ws", handleWebSocket) // Dashboard mux.HandleFunc("/api/routers", handleRouters) mux.HandleFunc("/api/events/list", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleListEvents)) mux.HandleFunc("/api/command", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleCommand)) mux.HandleFunc("/api/commands", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleListCommands)) mux.HandleFunc("/api/metrics", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleMetrics)) mux.HandleFunc("/api/alerts", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleAlerts)) mux.HandleFunc("/api/alerts/", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleAckAlert)) mux.HandleFunc("/api/auth/me", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleMe)) mux.HandleFunc("/api/events/stream", handleSSEStream) srv := &http.Server{ Addr: fmt.Sprintf(":%d", cfg.Port), Handler: mux, ReadTimeout: 0, WriteTimeout: 0, } go func() { sigCh := make(chan os.Signal, 1) signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM) <-sigCh log.Println("shutting down...") routersMu.RLock() for _, r := range routers { if r.Conn != nil { r.Conn.Close(websocket.StatusNormalClosure, "server shutdown") } } routersMu.RUnlock() shutdownCtx, c := context.WithTimeout(context.Background(), 5*time.Second) defer c() _ = srv.Shutdown(shutdownCtx) os.Exit(0) }() log.Printf("server ready on :%d", cfg.Port) if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) { log.Fatal(err) } } func pathBase(p string) string { i := strings.LastIndex(p, "/") if i < 0 { return p } return p[i+1:] } func ensureTopics(ctx context.Context) error { if kcl == nil { return nil } // Best-effort, log only log.Println("redpanda: topics will be auto-created on first publish") return nil }