main.go 9.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409
  1. // client2server - Go server with WebSocket + Redpanda
  2. // Copyright (c) 2026 Luis Rosales - MIT License
  3. //
  4. // Build: go build -o client2server-server
  5. // Run: ./client2server-server
  6. // WebSocket: ws://localhost:3843
  7. // HTTP API: http://localhost:3843
  8. package main
  9. import (
  10. "context"
  11. "encoding/json"
  12. "fmt"
  13. "log"
  14. "net/http"
  15. "os"
  16. "os/signal"
  17. "syscall"
  18. "time"
  19. "github.com/google/uuid"
  20. "github.com/redpanda-data/redpanda-sdk-go/redpanda"
  21. "github.com/redpanda-data/redpanda-sdk-go/schema"
  22. "nhooyr.io/websocket"
  23. )
  24. // Config
  25. type Config struct {
  26. RedpandaBrokers []string
  27. WebsocketPort int
  28. APIPort int
  29. Token string
  30. }
  31. var cfg = Config{
  32. RedpandaBrokers: []string{"localhost:9092"},
  33. WebsocketPort: 3843,
  34. APIPort: 3844, // HTTP API on next port
  35. Token: "secret-token",
  36. }
  37. // Event from router
  38. type RouterEvent struct {
  39. ID string `json:"id"`
  40. RouterID string `json:"router_id"`
  41. Hostname string `json:"hostname"`
  42. EventType string `json:"event_type"`
  43. Timestamp time.Time `json:"timestamp"`
  44. Payload map[string]interface{} `json:"payload"`
  45. ReceivedAt time.Time `json:"received_at"`
  46. Connection string `json:"connection"` // "websocket" or "http"
  47. }
  48. // Command to router
  49. type RouterCommand struct {
  50. ID string `json:"id"`
  51. RouterID string `json:"router_id"`
  52. Command string `json:"command"`
  53. Args map[string]string `json:"args,omitempty"`
  54. SentAt time.Time `json:"sent_at"`
  55. }
  56. // Router state
  57. type Router struct {
  58. ID string
  59. LastSeen time.Time
  60. Conn *websocket.Conn
  61. Connected bool
  62. }
  63. var routers = make(map[string]*Router)
  64. // Redpanda
  65. var rp *redpanda.Client
  66. func initRedpanda() error {
  67. cfg.RedpandaBrokers = getEnvComma("REDPANDA_BROKERS", "localhost:9092")
  68. var err error
  69. rp, err = redpanda.NewClient(&redpanda.ClientConfig{
  70. Brokers: cfg.RedpandaBrokers,
  71. })
  72. if err != nil {
  73. return fmt.Errorf("redpanda: %v", err)
  74. }
  75. // Create topics if needed
  76. topics := []string{"router-events", "router-commands"}
  77. for _, topic := range topics {
  78. err := rp.CreateTopic(topic, 1, 3)
  79. if err != nil && !schema.ErrTopicExists.Exists(err) {
  80. log.Printf("Topic %s: %v", topic, err)
  81. }
  82. }
  83. return nil
  84. }
  85. func getEnvComma(key, def string) []string {
  86. val := os.Getenv(key)
  87. if val == "" {
  88. return []string{def}
  89. }
  90. return []string{val}
  91. }
  92. // Publish event to Redpanda
  93. func publishEvent(event RouterEvent) error {
  94. data, err := json.Marshal(event)
  95. if err != nil {
  96. return err
  97. }
  98. return rp.Produce("router-events", []byte(event.ID), data)
  99. }
  100. // Publish command to Redpanda
  101. func publishCommand(cmd RouterCommand) error {
  102. data, err := json.Marshal(cmd)
  103. if err != nil {
  104. return err
  105. }
  106. return rp.Produce("router-commands", []byte(cmd.ID), data)
  107. }
  108. // WebSocket handler
  109. func handleWebSocket(w http.ResponseWriter, r *http.Request) {
  110. // Auth
  111. token := r.URL.Query().Get("token")
  112. if token != cfg.Token {
  113. http.Error(w, "Unauthorized", http.StatusUnauthorized)
  114. return
  115. }
  116. conn, err := websocket.Accept(w, r, &websocket.AcceptOptions{
  117. CompressionMode: websocket.CompressionContextTakeover,
  118. })
  119. if err != nil {
  120. log.Printf("WS accept: %v", err)
  121. return
  122. }
  123. defer conn.Close(websocket.StatusNormalClosure, "")
  124. ctx := context.Background()
  125. // Read router registration
  126. var regMsg RouterEvent
  127. err = conn.Read(ctx, &regMsg)
  128. if err != nil {
  129. log.Printf("WS read reg: %v", err)
  130. return
  131. }
  132. routerID := regMsg.RouterID
  133. if routerID == "" {
  134. routerID = r.RemoteAddr
  135. }
  136. // Store router
  137. routers[routerID] = &Router{
  138. ID: routerID,
  139. LastSeen: time.Now(),
  140. Conn: conn,
  141. Connected: true,
  142. }
  143. log.Printf("Router connected: %s", routerID)
  144. // Update router last seen
  145. routers[routerID].LastSeen = time.Now()
  146. // Message loop
  147. for {
  148. var event RouterEvent
  149. err := conn.Read(ctx, &event)
  150. if err != nil {
  151. break
  152. }
  153. event.ID = uuid.New().String()
  154. event.ReceivedAt = time.Now()
  155. event.Connection = "websocket"
  156. // Publish to Redpanda
  157. if err := publishEvent(event); err != nil {
  158. log.Printf("Publish error: %v", err)
  159. }
  160. log.Printf("[%s] %s", routerID, event.EventType)
  161. // Send ack
  162. conn.Write(ctx, []byte(`{"ack":true}`))
  163. // Update last seen
  164. routers[routerID].LastSeen = time.Now()
  165. }
  166. // Cleanup
  167. if routers[routerID] != nil {
  168. routers[routerID].Connected = false
  169. }
  170. log.Printf("Router disconnected: %s", routerID)
  171. }
  172. // HTTP Event webhook (fallback)
  173. func handleHTTPEvent(w http.ResponseWriter, r *http.Request) {
  174. if r.Method != "POST" {
  175. http.Error(w, "Method not allowed", http.StatusMethodNotAllowed)
  176. return
  177. }
  178. token := r.Header.Get("Authorization")
  179. if token != "Bearer "+cfg.Token && token != cfg.Token {
  180. http.Error(w, "Unauthorized", http.StatusUnauthorized)
  181. return
  182. }
  183. var event RouterEvent
  184. if err := json.NewDecoder(r.Body).Decode(&event); err != nil {
  185. http.Error(w, "Invalid JSON", http.StatusBadRequest)
  186. return
  187. }
  188. event.ID = uuid.New().String()
  189. event.ReceivedAt = time.Now()
  190. event.Connection = "http"
  191. // Publish to Redpanda
  192. if err := publishEvent(event); err != nil {
  193. log.Printf("Publish error: %v", err)
  194. w.WriteHeader(http.StatusInternalServerError)
  195. return
  196. }
  197. log.Printf("[%s] %s (HTTP)", event.RouterID, event.EventType)
  198. json.NewEncoder(w).Encode(map[string]string{"event_id": event.ID})
  199. }
  200. // Get routers list
  201. func handleRouters(w http.ResponseWriter, r *http.Request) {
  202. json.NewEncoder(w).Encode(map[string]interface{}{
  203. "routers": func() []map[string]interface{} {
  204. list := []map[string]interface{}{}
  205. for id, router := range routers {
  206. list = append(list, map[string]interface{}{
  207. "id": id,
  208. "last_seen": router.LastSeen,
  209. "online": time.Since(router.LastSeen) < 60*time.Second,
  210. })
  211. }
  212. return list
  213. }(),
  214. })
  215. }
  216. // Get events
  217. func handleEvents(w http.ResponseWriter, r *http.Request) {
  218. // Simple consumer - in production use proper offset management
  219. consumer, err := rp.NewConsumer("router-events", "client2server-group")
  220. if err != nil {
  221. http.Error(w, err.Error(), http.StatusInternalServerError)
  222. return
  223. }
  224. defer consumer.Close()
  225. limit := 100
  226. if l := r.URL.Query().Get("limit"); l != "" {
  227. fmt.Sscanf(l, "%d", &limit)
  228. }
  229. events := []RouterEvent{}
  230. ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
  231. defer cancel()
  232. for i := 0; i < limit; i++ {
  233. msg, err := consumer.Consume(ctx, 5*time.Second)
  234. if err != nil {
  235. break
  236. }
  237. var event RouterEvent
  238. if json.Unmarshal(msg.Value, &event) == nil {
  239. events = append(events, event)
  240. }
  241. }
  242. json.NewEncoder(w).Encode(map[string]interface{}{
  243. "events": events,
  244. "total": len(events),
  245. })
  246. }
  247. // Send command to router
  248. func handleCommand(w http.ResponseWriter, r *http.Request) {
  249. if r.Method != "POST" {
  250. http.Error(w, "Method not allowed", http.StatusMethodNotAllowed)
  251. return
  252. }
  253. var cmd RouterCommand
  254. if err := json.NewDecoder(r.Body).Decode(&cmd); err != nil {
  255. http.Error(w, "Invalid JSON", http.StatusBadRequest)
  256. return
  257. }
  258. routerID := r.URL.Query().Get("router_id")
  259. if routerID == "" {
  260. routerID = cmd.RouterID
  261. }
  262. if routerID == "" {
  263. http.Error(w, "router_id required", http.StatusBadRequest)
  264. return
  265. }
  266. // Publish command
  267. cmd.ID = uuid.New().String()
  268. cmd.RouterID = routerID
  269. cmd.SentAt = time.Now()
  270. if err := publishCommand(cmd); err != nil {
  271. log.Printf("Publish command: %v", err)
  272. http.Error(w, err.Error(), http.StatusInternalServerError)
  273. return
  274. }
  275. log.Printf("Command to %s: %s", routerID, cmd.Command)
  276. json.NewEncoder(w).Encode(map[string]string{
  277. "command_id": cmd.ID,
  278. "status": "sent",
  279. })
  280. }
  281. // Health check
  282. func handleHealth(w http.ResponseWriter, r *http.Request) {
  283. json.NewEncoder(w).Encode(map[string]interface{}{
  284. "status": "ok",
  285. "routers": len(routers),
  286. "timestamp": time.Now(),
  287. })
  288. }
  289. // Main
  290. func main() {
  291. // Config from env
  292. cfg.Token = getEnvStr("TOKEN", cfg.Token)
  293. cfg.WebsocketPort = getEnvInt("PORT", cfg.WebsocketPort)
  294. log.SetFlags(log.LstdFlags | log.Lshortfile)
  295. log.Printf("=== client2server Go Server ===")
  296. log.Printf("WebSocket: ws://localhost:%d", cfg.WebsocketPort)
  297. log.Printf("HTTP API: http://localhost:%d", cfg.APIPort)
  298. // Init Redpanda
  299. if err := initRedpanda(); err != nil {
  300. log.Printf("Redpanda init failed (running without): %v", err)
  301. } else {
  302. log.Printf("Redpanda: %v", cfg.RedpandaBrokers)
  303. }
  304. // HTTP handlers
  305. http.HandleFunc("/", handleHealth)
  306. http.HandleFunc("/health", handleHealth)
  307. http.HandleFunc("/api/events", handleHTTPEvent)
  308. http.HandleFunc("/api/events/", handleHTTPEvent)
  309. http.HandleFunc("/api/routers", handleRouters)
  310. http.HandleFunc("/api/command", handleCommand)
  311. // WebSocket (HAProxy protocol)
  312. http.HandleFunc("/ws", handleWebSocket)
  313. // Graceful shutdown
  314. go func() {
  315. sigCh := make(chan os.Signal, 1)
  316. signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
  317. <-sigCh
  318. log.Println("Shutting down...")
  319. for _, router := range routers {
  320. router.Conn.Close(websocket.StatusNormalClosure, "server shutdown")
  321. }
  322. os.Exit(0)
  323. }()
  324. // Start (using HTTP for both WS and API - need haproxy or prefix)
  325. log.Printf("Server ready on port %d", cfg.WebiscrollPoint)
  326. log.Fatal(http.ListenAndServe(fmt.Sprintf(":%d", cfg.WebsocketPort), nil))
  327. }
  328. func getEnvStr(key, def string) string {
  329. if val := os.Getenv(key); val != "" {
  330. return val
  331. }
  332. return def
  333. }
  334. func getEnvInt(key string, def int) int {
  335. if val := os.Getenv(key); val != "" {
  336. var v int
  337. fmt.Sscanf(val, "%d", &v)
  338. return v
  339. }
  340. return def
  341. }