diff --git a/serv.go b/serv.go new file mode 100644 index 0000000..89b138e --- /dev/null +++ b/serv.go @@ -0,0 +1,292 @@ +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) +}