package main import ( "context" "flag" "net/http" "os" "os/signal" "syscall" "time" "github.com/confluentinc/confluent-kafka-go/kafka" "github.com/gorilla/handlers" "github.com/gorilla/mux" "github.com/gorilla/websocket" uuid "github.com/satori/go.uuid" log "github.com/sirupsen/logrus" ) const ( // Time allowed to write a message to the peer. writeWait = 10 * time.Second // Time allowed to read the next pong message from the peer. pongWait = 60 * time.Second // Send pings to peer with this period. Must be less than pongWait. pingPeriod = (pongWait * 9) / 10 // Maximum message size allowed from peer. maxMessageSize = 512 ) var ( newline = []byte{'\n'} space = []byte{' '} addr = flag.String("addr", ":9000", "http service address") bootstrapServers = flag.String("bootstrap-servers", "localhost:9092", "kafka bootstrap servers") upgrader = websocket.Upgrader{ ReadBufferSize: 1024, WriteBufferSize: 1024, } ) // Client is a middleman between the websocket connection and the hub. type Client struct { hub *Hub // The websocket connection. conn *websocket.Conn // Buffered channel of outbound messages. send chan []byte } // readPump pumps messages from the websocket connection to the hub. // // The application runs readPump in a per-connection goroutine. The application // ensures that there is at most one reader on a connection by executing all // reads from this goroutine. func (c *Client) readPump(topic string) { defer func() { c.hub.unregister <- c c.conn.Close() }() c.conn.SetReadLimit(maxMessageSize) c.conn.SetReadDeadline(time.Now().Add(pongWait)) c.conn.SetPongHandler(func(string) error { c.conn.SetReadDeadline(time.Now().Add(pongWait)); return nil }) sigchan := make(chan os.Signal, 1) signal.Notify(sigchan, syscall.SIGINT, syscall.SIGTERM) group := uuid.Must(uuid.NewV4()) log.WithFields(log.Fields{ "topic": topic, "bootstrap servers": *bootstrapServers, "group": group}).Info("create kafka client") consumer, err := kafka.NewConsumer(&kafka.ConfigMap{ "bootstrap.servers": *bootstrapServers, "group.id": group, "session.timeout.ms": 6000, "auto.offset.reset": "latest"}) if err != nil { // fmt.Fprintf(os.Stderr, "Failed to create consumer: %s\n", err) log.WithError(err).Fatal("Failed to create consumer") os.Exit(1) } log.WithFields(log.Fields{ "Consumer": consumer, }).Info("Created ") consumer.SubscribeTopics([]string{topic}, nil) run := true for run == true { select { case sig := <-sigchan: log.WithFields(log.Fields{"signal": sig}).Info(" terminating") run = false log.Fatal("break ") default: ev := consumer.Poll(100) if ev == nil { continue } switch e := ev.(type) { case *kafka.Message: log.WithFields(log.Fields{"partition": e.TopicPartition, "message": string(e.Value)}).Info("got message") c.hub.broadcast <- e.Value if e.Headers != nil { log.WithFields(log.Fields{"header": e.Headers}).Info(" with header") } case kafka.Error: // Errors should generally be considered as informational, the client will try to automatically recover log.WithFields(log.Fields{"error": e}).Error("got error") default: log.WithFields(log.Fields{"event": e}).Info("ignored") } } } consumer.Close() } // writePump pumps messages from the hub to the websocket connection. // // A goroutine running writePump is started for each connection. The // application ensures that there is at most one writer to a connection by // executing all writes from this goroutine. func (c *Client) writePump() { ticker := time.NewTicker(pingPeriod) defer func() { ticker.Stop() c.conn.Close() }() for { select { case message, ok := <-c.send: c.conn.SetWriteDeadline(time.Now().Add(writeWait)) if !ok { // The hub closed the channel. c.conn.WriteMessage(websocket.CloseMessage, []byte{}) return } w, err := c.conn.NextWriter(websocket.TextMessage) if err != nil { return } w.Write(message) // Add queued chat messages to the current websocket message. n := len(c.send) for i := 0; i < n; i++ { w.Write(newline) w.Write(<-c.send) } if err := w.Close(); err != nil { return } case <-ticker.C: c.conn.SetWriteDeadline(time.Now().Add(writeWait)) if err := c.conn.WriteMessage(websocket.PingMessage, nil); err != nil { return } } } } // serveWs handles websocket requests from the peer. func serveWs(hub *Hub, w http.ResponseWriter, r *http.Request) { vars := mux.Vars(r) topic := vars["topic"] conn, err := upgrader.Upgrade(w, r, nil) if err != nil { log.Error("upgrade connection failed ", err) return } client := &Client{hub: hub, conn: conn, send: make(chan []byte, 256)} client.hub.register <- client // Allow collection of memory referenced by the caller by doing all work in // new goroutines. go client.writePump() go client.readPump(topic) } // Hub maintains the set of active clients and broadcasts messages to the // clients. type Hub struct { // Registered clients. clients map[*Client]bool // Inbound messages from the clients. broadcast chan []byte // Register requests from the clients. register chan *Client // Unregister requests from clients. unregister chan *Client } func newHub() *Hub { return &Hub{ broadcast: make(chan []byte), register: make(chan *Client), unregister: make(chan *Client), clients: make(map[*Client]bool), } } func (h *Hub) run() { for { select { case client := <-h.register: h.clients[client] = true case client := <-h.unregister: if _, ok := h.clients[client]; ok { delete(h.clients, client) close(client.send) } case message := <-h.broadcast: for client := range h.clients { select { case client.send <- message: default: close(client.send) delete(h.clients, client) } } } } } func main() { var wait time.Duration flag.DurationVar(&wait, "graceful-timeout", time.Second*15, "the duration for which the server gracefully wait for existing connections to finish - e.g. 15s or 1m") flag.Parse() r := mux.NewRouter() loggedRouter := handlers.LoggingHandler(os.Stdout, r) srv := &http.Server{ Addr: *addr, // Good practice to set timeouts to avoid Slowloris attacks. WriteTimeout: time.Second * 15, ReadTimeout: time.Second * 15, IdleTimeout: time.Second * 60, Handler: loggedRouter, // Pass our instance of gorilla/mux in. } hub := newHub() go hub.run() r.HandleFunc("/ws/{topic}", func(w http.ResponseWriter, r *http.Request) { serveWs(hub, w, r) }) // Run our server in a goroutine so that it doesn't block. go func() { if err := srv.ListenAndServe(); err != nil { log.WithFields(log.Fields{"error": err}).Error("startup error") } else { log.Info("proxy started at :: ", addr) } }() c := make(chan os.Signal, 1) // We'll accept graceful shutdowns when quit via SIGINT (Ctrl+C) // SIGKILL, SIGQUIT or SIGTERM (Ctrl+/) will not be caught. signal.Notify(c, os.Interrupt) // Block until we receive our signal. <-c // Create a deadline to wait for. ctx, cancel := context.WithTimeout(context.Background(), wait) defer cancel() // Doesn't block if no connections, but will otherwise wait // until the timeout deadline. srv.Shutdown(ctx) // Optionally, you could run srv.Shutdown in a goroutine and block on // <-ctx.Done() if your application should wait for other services // to finalize based on context cancellation. log.Info("shutting down") os.Exit(0) }