| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892 |
- // 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
- }
|