2019-01-05 16:06:47 +08:00
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 () {
2019-01-05 16:13:16 +08:00
log . WithFields ( log . Fields { "address" : * addr }). Info ( "starting server" )
2019-01-05 16:06:47 +08:00
if err := srv . ListenAndServe (); err != nil {
log . WithFields ( log . Fields { "error" : err }). Error ( "startup error" )
}
}()
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 )
}