main.go 24 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892
  1. // client2server - Go server with WebSocket + Redpanda + Dashboard API
  2. // Copyright (c) 2026 Luis Rosales - MIT License
  3. //
  4. // Stack:
  5. // - WebSocket: github.com/coder/websocket
  6. // - Kafka: github.com/twmb/franz-go (talks to Redpanda)
  7. // - Auth: JWT (HS256) with role-based access
  8. // - Storage: SQLite (modernc.org/sqlite, pure Go, no cgo)
  9. // - Live feed: Server-Sent Events
  10. //
  11. // Endpoints:
  12. // POST /api/auth/login - username/password -> JWT
  13. // GET /api/auth/me - current user
  14. // GET /health - liveness
  15. // GET /api/routers - known routers (legacy shared token OK)
  16. // GET /api/events?limit=&router_id= - event history (auth)
  17. // POST /api/events - ingest event (router or hotplug)
  18. // POST /api/command - send command to router (auth)
  19. // GET /api/commands?router_id= - command history
  20. // GET /api/metrics?since=1h - time-series metrics
  21. // GET /api/alerts?unack=1 - alerts
  22. // POST /api/alerts/:id/ack - acknowledge alert
  23. // GET /api/events/stream - SSE live feed
  24. // GET /ws - WebSocket from routers (legacy token)
  25. //
  26. // Env: REDPANDA_BROKERS, TOKEN (legacy shared), JWT_SECRET, DB_PATH, PORT
  27. package main
  28. import (
  29. "context"
  30. "encoding/json"
  31. "errors"
  32. "fmt"
  33. "log"
  34. "net/http"
  35. "os"
  36. "os/signal"
  37. "strconv"
  38. "strings"
  39. "sync"
  40. "syscall"
  41. "time"
  42. "github.com/coder/websocket"
  43. "github.com/google/uuid"
  44. "github.com/twmb/franz-go/pkg/kgo"
  45. )
  46. // ----------------------------------------------------------------------------
  47. // Config
  48. // ----------------------------------------------------------------------------
  49. type Config struct {
  50. RedpandaBrokers []string
  51. Port int
  52. Token string // legacy shared token for routers/hotplug
  53. DBPath string
  54. }
  55. var cfg = Config{
  56. RedpandaBrokers: []string{"localhost:9092"},
  57. Port: 3843,
  58. Token: "secret-token",
  59. DBPath: "data/client2server.db",
  60. }
  61. func loadConfigFromEnv() {
  62. cfg.RedpandaBrokers = strings.Split(getEnvStr("REDPANDA_BROKERS", "localhost:9092"), ",")
  63. cfg.Port = getEnvInt("PORT", cfg.Port)
  64. cfg.Token = getEnvStr("TOKEN", cfg.Token)
  65. cfg.DBPath = getEnvStr("DB_PATH", cfg.DBPath)
  66. }
  67. func getEnvStr(key, def string) string {
  68. if v := os.Getenv(key); v != "" {
  69. return v
  70. }
  71. return def
  72. }
  73. func getEnvInt(key string, def int) int {
  74. if v := os.Getenv(key); v != "" {
  75. var n int
  76. if _, err := fmt.Sscanf(v, "%d", &n); err == nil {
  77. return n
  78. }
  79. }
  80. return def
  81. }
  82. // ----------------------------------------------------------------------------
  83. // Domain types
  84. // ----------------------------------------------------------------------------
  85. type RouterEvent struct {
  86. ID string `json:"id"`
  87. RouterID string `json:"router_id"`
  88. Hostname string `json:"hostname,omitempty"`
  89. EventType string `json:"event_type"`
  90. Timestamp time.Time `json:"timestamp"`
  91. Payload map[string]interface{} `json:"payload"`
  92. ReceivedAt time.Time `json:"received_at"`
  93. Connection string `json:"connection"`
  94. }
  95. type RouterCommand struct {
  96. ID string `json:"id"`
  97. RouterID string `json:"router_id"`
  98. Command string `json:"command"`
  99. Args map[string]string `json:"args,omitempty"`
  100. SentAt time.Time `json:"sent_at"`
  101. }
  102. type CommandResult struct {
  103. Success bool `json:"success"`
  104. Output string `json:"output"`
  105. Error string `json:"error"`
  106. }
  107. // ----------------------------------------------------------------------------
  108. // Router registry
  109. // ----------------------------------------------------------------------------
  110. type Router struct {
  111. ID string
  112. LastSeen time.Time
  113. Conn *websocket.Conn
  114. Connected bool
  115. writeMu sync.Mutex
  116. }
  117. var (
  118. routersMu sync.RWMutex
  119. routers = make(map[string]*Router)
  120. routerQueuesMu sync.Mutex
  121. routerQueues = make(map[string][]RouterCommand)
  122. pendingMu sync.Mutex
  123. pendingCmds = make(map[string]chan CommandResult)
  124. executedMu sync.Mutex
  125. executedCmds = make(map[string]time.Time)
  126. )
  127. const (
  128. commandTimeout = 30 * time.Second
  129. idempotencyTTL = 5 * time.Minute
  130. wsWriteTimeout = 10 * time.Second
  131. publishTimeout = 3 * time.Second
  132. offlineThreshold = 60 * time.Second
  133. )
  134. // ----------------------------------------------------------------------------
  135. // Redpanda (Kafka) producer
  136. // ----------------------------------------------------------------------------
  137. var kcl *kgo.Client
  138. func initRedpanda(ctx context.Context) error {
  139. cl, err := kgo.NewClient(
  140. kgo.SeedBrokers(cfg.RedpandaBrokers...),
  141. kgo.ClientID("client2server"),
  142. kgo.ProducerLinger(5*time.Millisecond),
  143. kgo.ProducerBatchCompression(kgo.SnappyCompression()),
  144. )
  145. if err != nil {
  146. return fmt.Errorf("kafka client: %w", err)
  147. }
  148. kcl = cl
  149. return nil
  150. }
  151. func publish(ctx context.Context, topic string, key string, value any) error {
  152. data, err := json.Marshal(value)
  153. if err != nil {
  154. return fmt.Errorf("marshal: %w", err)
  155. }
  156. rec := &kgo.Record{Topic: topic, Key: []byte(key), Value: data}
  157. // Fire-and-forget: don't block the HTTP path. franz-go queues records
  158. // internally and flushes them in the background. If the broker is down,
  159. // the record will be retried until it expires (RecordDeliveryTimeout).
  160. // We still get error visibility via a callback for logging.
  161. pctx, cancel := context.WithTimeout(ctx, 100*time.Millisecond)
  162. defer cancel()
  163. res := kcl.ProduceSync(pctx, rec)
  164. if err := res.FirstErr(); err != nil {
  165. log.Printf("publish %s/%s: %v", topic, key, err)
  166. return err
  167. }
  168. return nil
  169. }
  170. // ----------------------------------------------------------------------------
  171. // Event ingestion (called by both WS and HTTP)
  172. // ----------------------------------------------------------------------------
  173. func ingestEvent(ev RouterEvent) {
  174. if ev.ID == "" {
  175. ev.ID = uuid.New().String()
  176. }
  177. if ev.ReceivedAt.IsZero() {
  178. ev.ReceivedAt = time.Now()
  179. }
  180. if ev.Timestamp.IsZero() {
  181. ev.Timestamp = ev.ReceivedAt
  182. }
  183. // Metrics
  184. globalMetrics.recordEvent(ev.EventType)
  185. // Persist (best-effort)
  186. if err := saveEvent(ev); err != nil {
  187. log.Printf("save event: %v", err)
  188. }
  189. // Publish to Redpanda
  190. if err := publish(context.Background(), "router-events", ev.RouterID, ev); err != nil {
  191. log.Printf("publish event: %v", err)
  192. }
  193. // Broadcast to SSE clients
  194. hub.broadcast("event", ev)
  195. }
  196. // ----------------------------------------------------------------------------
  197. // WebSocket handler (routers connect here)
  198. // ----------------------------------------------------------------------------
  199. func handleWebSocket(w http.ResponseWriter, r *http.Request) {
  200. // Auth: legacy TOKEN only (routers use shared secret, not JWT)
  201. token := r.URL.Query().Get("token")
  202. if token == "" {
  203. if h := r.Header.Get("Authorization"); strings.HasPrefix(h, "Bearer ") {
  204. token = strings.TrimPrefix(h, "Bearer ")
  205. }
  206. }
  207. if token != cfg.Token {
  208. http.Error(w, "Unauthorized", http.StatusUnauthorized)
  209. return
  210. }
  211. conn, err := websocket.Accept(w, r, &websocket.AcceptOptions{
  212. CompressionMode: websocket.CompressionDisabled,
  213. })
  214. if err != nil {
  215. log.Printf("ws accept: %v", err)
  216. return
  217. }
  218. defer conn.Close(websocket.StatusNormalClosure, "bye")
  219. ctx, cancel := context.WithCancel(r.Context())
  220. defer cancel()
  221. // First message is registration
  222. raw, err := readRouterMessage(ctx, conn)
  223. if err != nil {
  224. log.Printf("ws read reg: %v", err)
  225. return
  226. }
  227. var reg RouterEvent
  228. if err := json.Unmarshal(raw, &reg); err != nil {
  229. log.Printf("ws reg parse: %v", err)
  230. return
  231. }
  232. routerID := reg.RouterID
  233. if routerID == "" {
  234. routerID = r.RemoteAddr
  235. }
  236. routersMu.Lock()
  237. routers[routerID] = &Router{
  238. ID: routerID,
  239. LastSeen: time.Now(),
  240. Conn: conn,
  241. Connected: true,
  242. }
  243. routersMu.Unlock()
  244. // Mark as back online (clears any offline alert)
  245. acknowledgeRouterAlerts(routerID)
  246. log.Printf("router connected: %s (from %s)", routerID, r.RemoteAddr)
  247. flushQueuedCommands(ctx, routerID)
  248. for {
  249. raw, err := readRouterMessage(ctx, conn)
  250. if err != nil {
  251. break
  252. }
  253. var ev RouterEvent
  254. if err := json.Unmarshal(raw, &ev); err != nil {
  255. log.Printf("[%s] bad event json: %v", routerID, err)
  256. continue
  257. }
  258. // Inherit router_id if missing
  259. if ev.RouterID == "" {
  260. ev.RouterID = routerID
  261. }
  262. ev.Connection = "websocket"
  263. // Command result handling
  264. if ev.EventType == "command_result" {
  265. if cid, _ := ev.Payload["command_id"].(string); cid != "" {
  266. deliverCommandResult(cid, ev.Payload)
  267. executedMu.Lock()
  268. executedCmds[cid] = time.Now()
  269. executedMu.Unlock()
  270. }
  271. }
  272. ingestEvent(ev)
  273. log.Printf("[%s] %s", routerID, ev.EventType)
  274. routersMu.Lock()
  275. if r, ok := routers[routerID]; ok {
  276. r.LastSeen = time.Now()
  277. }
  278. routersMu.Unlock()
  279. }
  280. routersMu.Lock()
  281. if r, ok := routers[routerID]; ok {
  282. r.Connected = false
  283. }
  284. routersMu.Unlock()
  285. log.Printf("router disconnected: %s", routerID)
  286. }
  287. func readRouterMessage(ctx context.Context, conn *websocket.Conn) ([]byte, error) {
  288. _, data, err := conn.Read(ctx)
  289. return data, err
  290. }
  291. func deliverCommandResult(cmdID string, payload map[string]interface{}) {
  292. pendingMu.Lock()
  293. ch, ok := pendingCmds[cmdID]
  294. if ok {
  295. delete(pendingCmds, cmdID)
  296. }
  297. pendingMu.Unlock()
  298. if !ok {
  299. return
  300. }
  301. res := CommandResult{Success: payload["success"] == true}
  302. if s, ok := payload["output"].(string); ok {
  303. res.Output = s
  304. }
  305. if s, ok := payload["error"].(string); ok {
  306. res.Error = s
  307. }
  308. // Persist result to SQLite
  309. _ = updateCommandResult(cmdID, "completed", &res)
  310. globalMetrics.recordCommand(res.Success)
  311. select {
  312. case ch <- res:
  313. default:
  314. }
  315. }
  316. func flushQueuedCommands(ctx context.Context, routerID string) {
  317. routerQueuesMu.Lock()
  318. queue := routerQueues[routerID]
  319. delete(routerQueues, routerID)
  320. routerQueuesMu.Unlock()
  321. if len(queue) == 0 {
  322. return
  323. }
  324. log.Printf("flushing %d queued commands to %s", len(queue), routerID)
  325. for _, cmd := range queue {
  326. cmd.SentAt = time.Now()
  327. pendingMu.Lock()
  328. pendingCmds[cmd.ID] = make(chan CommandResult, 1)
  329. pendingMu.Unlock()
  330. _ = updateCommandResult(cmd.ID, "delivered", nil)
  331. if err := publish(ctx, "router-commands", routerID, cmd); err != nil {
  332. log.Printf("queue flush publish: %v", err)
  333. }
  334. }
  335. }
  336. func acknowledgeRouterAlerts(routerID string) {
  337. // Best-effort: mark all unacked offline alerts for this router as resolved
  338. if db == nil {
  339. return
  340. }
  341. _, _ = db.Exec("UPDATE alerts SET acknowledged_at = CURRENT_TIMESTAMP WHERE router_id = ? AND kind = 'router_offline' AND acknowledged_at IS NULL", routerID)
  342. }
  343. // ----------------------------------------------------------------------------
  344. // HTTP API
  345. // ----------------------------------------------------------------------------
  346. func handleHealth(w http.ResponseWriter, r *http.Request) {
  347. routersMu.RLock()
  348. onlineCount := 0
  349. total := len(routers)
  350. for _, rt := range routers {
  351. if time.Since(rt.LastSeen) < offlineThreshold {
  352. onlineCount++
  353. }
  354. }
  355. routersMu.RUnlock()
  356. w.Header().Set("Content-Type", "application/json")
  357. _ = json.NewEncoder(w).Encode(map[string]interface{}{
  358. "status": "ok",
  359. "routers": total,
  360. "routers_online": onlineCount,
  361. "redpanda": cfg.RedpandaBrokers,
  362. "version": "2.1.0",
  363. })
  364. }
  365. func handleRouters(w http.ResponseWriter, r *http.Request) {
  366. type entry struct {
  367. ID string `json:"id"`
  368. LastSeen time.Time `json:"last_seen"`
  369. Online bool `json:"online"`
  370. Queued int `json:"queued_commands"`
  371. }
  372. routersMu.RLock()
  373. list := make([]entry, 0, len(routers))
  374. for id, rt := range routers {
  375. list = append(list, entry{
  376. ID: id,
  377. LastSeen: rt.LastSeen,
  378. Online: time.Since(rt.LastSeen) < offlineThreshold,
  379. Queued: len(routerQueues[id]),
  380. })
  381. }
  382. routersMu.RUnlock()
  383. w.Header().Set("Content-Type", "application/json")
  384. _ = json.NewEncoder(w).Encode(map[string]interface{}{"routers": list})
  385. }
  386. func handleHTTPEvent(w http.ResponseWriter, r *http.Request) {
  387. if r.Method != http.MethodPost {
  388. http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
  389. return
  390. }
  391. // Accept either JWT or legacy shared TOKEN
  392. if !authorisedRouter(r) {
  393. http.Error(w, "unauthorized", http.StatusUnauthorized)
  394. return
  395. }
  396. var ev RouterEvent
  397. if err := json.NewDecoder(r.Body).Decode(&ev); err != nil {
  398. http.Error(w, "invalid json", http.StatusBadRequest)
  399. return
  400. }
  401. ev.Connection = "http"
  402. ingestEvent(ev)
  403. w.Header().Set("Content-Type", "application/json")
  404. _ = json.NewEncoder(w).Encode(map[string]string{
  405. "event_id": ev.ID,
  406. "router_id": ev.RouterID,
  407. "status": "accepted",
  408. })
  409. _ = ev
  410. }
  411. func handleListEvents(w http.ResponseWriter, r *http.Request) {
  412. limit, _ := strconv.Atoi(r.URL.Query().Get("limit"))
  413. if limit <= 0 || limit > 1000 {
  414. limit = 100
  415. }
  416. routerID := r.URL.Query().Get("router_id")
  417. eventType := r.URL.Query().Get("event_type")
  418. events, err := listEvents(limit, routerID, eventType)
  419. if err != nil {
  420. http.Error(w, "db error: "+err.Error(), http.StatusInternalServerError)
  421. return
  422. }
  423. w.Header().Set("Content-Type", "application/json")
  424. _ = json.NewEncoder(w).Encode(map[string]interface{}{"events": events})
  425. }
  426. func handleCommand(w http.ResponseWriter, r *http.Request) {
  427. if r.Method != http.MethodPost {
  428. http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
  429. return
  430. }
  431. // JWT or legacy TOKEN
  432. if !authorisedAny(r) {
  433. http.Error(w, "unauthorized", http.StatusUnauthorized)
  434. return
  435. }
  436. var cmd RouterCommand
  437. if err := json.NewDecoder(r.Body).Decode(&cmd); err != nil {
  438. http.Error(w, "invalid json", http.StatusBadRequest)
  439. return
  440. }
  441. routerID := r.URL.Query().Get("router_id")
  442. if routerID == "" {
  443. routerID = cmd.RouterID
  444. }
  445. if routerID == "" {
  446. http.Error(w, "router_id required", http.StatusBadRequest)
  447. return
  448. }
  449. // Idempotency
  450. xecID := cmd.ID
  451. if xecID == "" {
  452. xecID = uuid.New().String()
  453. }
  454. executedMu.Lock()
  455. if last, exists := executedCmds[xecID]; exists && time.Since(last) < idempotencyTTL {
  456. executedMu.Unlock()
  457. w.Header().Set("Content-Type", "application/json")
  458. _ = json.NewEncoder(w).Encode(map[string]interface{}{
  459. "idempotent_reject": true,
  460. "existing_command_id": xecID,
  461. "message": "command already executed within TTL",
  462. })
  463. return
  464. }
  465. delete(executedCmds, xecID)
  466. executedMu.Unlock()
  467. cmd.ID = xecID
  468. cmd.RouterID = routerID
  469. cmd.SentAt = time.Now()
  470. // Identify the issuer (JWT username) if present
  471. issuedBy := "token"
  472. if claims := claimsFromHeader(r); claims != nil {
  473. issuedBy = claims.Username
  474. }
  475. // Persist
  476. _ = saveCommand(cmd, issuedBy, "pending")
  477. routersMu.RLock()
  478. router, online := routers[routerID]
  479. routersMu.RUnlock()
  480. if !online || router == nil || !router.Connected {
  481. routerQueuesMu.Lock()
  482. routerQueues[routerID] = append(routerQueues[routerID], cmd)
  483. depth := len(routerQueues[routerID])
  484. routerQueuesMu.Unlock()
  485. log.Printf("queued cmd %s for %s (depth=%d)", cmd.Command, routerID, depth)
  486. _ = updateCommandResult(cmd.ID, "queued", nil)
  487. w.Header().Set("Content-Type", "application/json")
  488. _ = json.NewEncoder(w).Encode(map[string]interface{}{
  489. "command_id": cmd.ID,
  490. "status": "queued",
  491. "queued_for": routerID,
  492. "queue_depth": depth,
  493. })
  494. return
  495. }
  496. resultCh := make(chan CommandResult, 1)
  497. pendingMu.Lock()
  498. pendingCmds[cmd.ID] = resultCh
  499. pendingMu.Unlock()
  500. if err := publish(r.Context(), "router-commands", routerID, cmd); err != nil {
  501. pendingMu.Lock()
  502. delete(pendingCmds, cmd.ID)
  503. pendingMu.Unlock()
  504. _ = updateCommandResult(cmd.ID, "publish_failed", nil)
  505. http.Error(w, "publish failed: "+err.Error(), http.StatusBadGateway)
  506. return
  507. }
  508. _ = updateCommandResult(cmd.ID, "delivered", nil)
  509. log.Printf("cmd %s -> %s (waiting)", cmd.Command, routerID)
  510. w.Header().Set("Content-Type", "application/json")
  511. select {
  512. case res := <-resultCh:
  513. _ = json.NewEncoder(w).Encode(map[string]interface{}{
  514. "command_id": cmd.ID,
  515. "status": "completed",
  516. "success": res.Success,
  517. "output": res.Output,
  518. "error": res.Error,
  519. })
  520. case <-time.After(commandTimeout):
  521. pendingMu.Lock()
  522. delete(pendingCmds, cmd.ID)
  523. pendingMu.Unlock()
  524. _ = updateCommandResult(cmd.ID, "timeout", nil)
  525. _ = json.NewEncoder(w).Encode(map[string]string{
  526. "command_id": cmd.ID,
  527. "status": "timeout",
  528. "error": "router did not respond",
  529. })
  530. }
  531. }
  532. func handleListCommands(w http.ResponseWriter, r *http.Request) {
  533. limit, _ := strconv.Atoi(r.URL.Query().Get("limit"))
  534. if limit <= 0 || limit > 500 {
  535. limit = 50
  536. }
  537. routerID := r.URL.Query().Get("router_id")
  538. cmds, err := listCommands(limit, routerID)
  539. if err != nil {
  540. http.Error(w, "db error: "+err.Error(), http.StatusInternalServerError)
  541. return
  542. }
  543. w.Header().Set("Content-Type", "application/json")
  544. _ = json.NewEncoder(w).Encode(map[string]interface{}{"commands": cmds})
  545. }
  546. func handleMetrics(w http.ResponseWriter, r *http.Request) {
  547. sinceStr := r.URL.Query().Get("since")
  548. since := 1 * time.Hour
  549. if sinceStr != "" {
  550. if d, err := time.ParseDuration(sinceStr); err == nil {
  551. since = d
  552. }
  553. }
  554. w.Header().Set("Content-Type", "application/json")
  555. _ = json.NewEncoder(w).Encode(map[string]interface{}{
  556. "summary": globalMetrics.summary(),
  557. "buckets": globalMetrics.snapshot(since),
  558. "since": since.String(),
  559. })
  560. }
  561. func handleAlerts(w http.ResponseWriter, r *http.Request) {
  562. if r.Method == http.MethodPost {
  563. // Acknowledge: POST /api/alerts/{id}/ack
  564. // Path is set up by the mux (see main)
  565. http.Error(w, "use POST /api/alerts/{id}/ack", http.StatusMethodNotAllowed)
  566. return
  567. }
  568. unack := r.URL.Query().Get("unack") == "1"
  569. limit, _ := strconv.Atoi(r.URL.Query().Get("limit"))
  570. if limit <= 0 || limit > 500 {
  571. limit = 100
  572. }
  573. alerts, err := listAlerts(unack, limit)
  574. if err != nil {
  575. http.Error(w, "db error: "+err.Error(), http.StatusInternalServerError)
  576. return
  577. }
  578. w.Header().Set("Content-Type", "application/json")
  579. _ = json.NewEncoder(w).Encode(map[string]interface{}{"alerts": alerts})
  580. }
  581. func handleAckAlert(w http.ResponseWriter, r *http.Request) {
  582. if r.Method != http.MethodPost {
  583. http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
  584. return
  585. }
  586. // Extract id from /api/alerts/{id}/ack
  587. parts := strings.Split(strings.Trim(r.URL.Path, "/"), "/")
  588. if len(parts) < 3 {
  589. http.Error(w, "bad path", http.StatusBadRequest)
  590. return
  591. }
  592. id, err := strconv.ParseInt(parts[2], 10, 64)
  593. if err != nil {
  594. http.Error(w, "bad id", http.StatusBadRequest)
  595. return
  596. }
  597. if err := acknowledgeAlert(id); err != nil {
  598. http.Error(w, err.Error(), http.StatusInternalServerError)
  599. return
  600. }
  601. w.Header().Set("Content-Type", "application/json")
  602. _ = json.NewEncoder(w).Encode(map[string]string{"status": "acknowledged"})
  603. }
  604. func handleMe(w http.ResponseWriter, r *http.Request) {
  605. claims := claimsFromHeader(r)
  606. if claims == nil {
  607. http.Error(w, "unauthorized", http.StatusUnauthorized)
  608. return
  609. }
  610. w.Header().Set("Content-Type", "application/json")
  611. _ = json.NewEncoder(w).Encode(map[string]interface{}{
  612. "user_id": claims.UserID,
  613. "username": claims.Username,
  614. "role": claims.Role,
  615. "expires": claims.ExpiresAt,
  616. })
  617. }
  618. // ----------------------------------------------------------------------------
  619. // Authorisation helpers
  620. // ----------------------------------------------------------------------------
  621. func claimsFromHeader(r *http.Request) *Claims {
  622. auth := r.Header.Get("Authorization")
  623. token := trimBearer(auth)
  624. if token == "" || token == cfg.Token {
  625. return nil
  626. }
  627. claims, err := ParseJWT(token)
  628. if err != nil {
  629. return nil
  630. }
  631. return claims
  632. }
  633. func authorisedAny(r *http.Request) bool {
  634. auth := r.Header.Get("Authorization")
  635. token := trimBearer(auth)
  636. if token == "" {
  637. return false
  638. }
  639. if token == cfg.Token {
  640. return true
  641. }
  642. _, err := ParseJWT(token)
  643. return err == nil
  644. }
  645. func authorisedRouter(r *http.Request) bool {
  646. auth := r.Header.Get("Authorization")
  647. token := trimBearer(auth)
  648. if token == "" {
  649. return false
  650. }
  651. // Routers use the legacy shared TOKEN. JWT users are valid too (in case
  652. // someone scripts an event submission from the dashboard).
  653. if token == cfg.Token {
  654. return true
  655. }
  656. _, err := ParseJWT(token)
  657. return err == nil
  658. }
  659. // ----------------------------------------------------------------------------
  660. // Background jobs
  661. // ----------------------------------------------------------------------------
  662. func startJanitor(ctx context.Context) {
  663. go func() {
  664. t := time.NewTicker(time.Minute)
  665. defer t.Stop()
  666. for {
  667. select {
  668. case <-ctx.Done():
  669. return
  670. case <-t.C:
  671. now := time.Now()
  672. executedMu.Lock()
  673. for id, ts := range executedCmds {
  674. if now.Sub(ts) > idempotencyTTL {
  675. delete(executedCmds, id)
  676. }
  677. }
  678. executedMu.Unlock()
  679. }
  680. }
  681. }()
  682. }
  683. func startOfflineWatcher(ctx context.Context) {
  684. go func() {
  685. t := time.NewTicker(30 * time.Second)
  686. defer t.Stop()
  687. known := map[string]bool{}
  688. for {
  689. select {
  690. case <-ctx.Done():
  691. return
  692. case <-t.C:
  693. routersMu.RLock()
  694. current := map[string]bool{}
  695. for id, rt := range routers {
  696. online := time.Since(rt.LastSeen) < offlineThreshold
  697. current[id] = online
  698. if !online && !known[id] {
  699. createAlert(id, "router_offline", fmt.Sprintf("Router %s has been offline for %s", id, time.Since(rt.LastSeen).Round(time.Second)))
  700. }
  701. }
  702. routersMu.RUnlock()
  703. known = current
  704. }
  705. }
  706. }()
  707. }
  708. // ----------------------------------------------------------------------------
  709. // main
  710. // ----------------------------------------------------------------------------
  711. func main() {
  712. loadConfigFromEnv()
  713. log.SetFlags(log.LstdFlags | log.Lshortfile)
  714. log.Printf("=== client2server v2.1 ===")
  715. log.Printf("port=%d brokers=%v db=%s", cfg.Port, cfg.RedpandaBrokers, cfg.DBPath)
  716. // Ensure DB directory exists
  717. if err := os.MkdirAll(strings.TrimSuffix(cfg.DBPath, "/"+pathBase(cfg.DBPath)), 0755); err != nil {
  718. log.Printf("mkdir db: %v", err)
  719. }
  720. if err := initStore(cfg.DBPath); err != nil {
  721. log.Fatalf("init store: %v", err)
  722. }
  723. ctx, cancel := context.WithCancel(context.Background())
  724. defer cancel()
  725. if err := initRedpanda(ctx); err != nil {
  726. log.Printf("redpanda init failed (continuing): %v", err)
  727. } else {
  728. defer kcl.Close()
  729. // Try to ensure topics exist (best effort)
  730. if err := ensureTopics(ctx); err != nil {
  731. log.Printf("ensure topics: %v", err)
  732. }
  733. // Start the consumer that persists to SQLite
  734. go startConsumer(ctx)
  735. }
  736. startJanitor(ctx)
  737. startOfflineWatcher(ctx)
  738. mux := http.NewServeMux()
  739. // Public
  740. mux.HandleFunc("/health", handleHealth)
  741. mux.HandleFunc("/", handleHealth)
  742. // Auth
  743. mux.HandleFunc("/api/auth/login", handleLogin)
  744. // Router-facing (legacy token OR JWT)
  745. mux.HandleFunc("/api/events", handleHTTPEvent)
  746. mux.HandleFunc("/api/events/", handleHTTPEvent)
  747. mux.HandleFunc("/ws", handleWebSocket)
  748. // Dashboard
  749. mux.HandleFunc("/api/routers", handleRouters)
  750. mux.HandleFunc("/api/events/list", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleListEvents))
  751. mux.HandleFunc("/api/command", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleCommand))
  752. mux.HandleFunc("/api/commands", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleListCommands))
  753. mux.HandleFunc("/api/metrics", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleMetrics))
  754. mux.HandleFunc("/api/alerts", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleAlerts))
  755. mux.HandleFunc("/api/alerts/", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleAckAlert))
  756. mux.HandleFunc("/api/auth/me", requireRole(RoleUser, RoleProjectAdmin, RoleSystemAdmin)(handleMe))
  757. mux.HandleFunc("/api/events/stream", handleSSEStream)
  758. srv := &http.Server{
  759. Addr: fmt.Sprintf(":%d", cfg.Port),
  760. Handler: mux,
  761. ReadTimeout: 0,
  762. WriteTimeout: 0,
  763. }
  764. go func() {
  765. sigCh := make(chan os.Signal, 1)
  766. signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
  767. <-sigCh
  768. log.Println("shutting down...")
  769. routersMu.RLock()
  770. for _, r := range routers {
  771. if r.Conn != nil {
  772. r.Conn.Close(websocket.StatusNormalClosure, "server shutdown")
  773. }
  774. }
  775. routersMu.RUnlock()
  776. shutdownCtx, c := context.WithTimeout(context.Background(), 5*time.Second)
  777. defer c()
  778. _ = srv.Shutdown(shutdownCtx)
  779. os.Exit(0)
  780. }()
  781. log.Printf("server ready on :%d", cfg.Port)
  782. if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
  783. log.Fatal(err)
  784. }
  785. }
  786. func pathBase(p string) string {
  787. i := strings.LastIndex(p, "/")
  788. if i < 0 {
  789. return p
  790. }
  791. return p[i+1:]
  792. }
  793. func ensureTopics(ctx context.Context) error {
  794. if kcl == nil {
  795. return nil
  796. }
  797. // Best-effort, log only
  798. log.Println("redpanda: topics will be auto-created on first publish")
  799. return nil
  800. }