2019-01-05 16:06:47 +08:00
package main
import (
"context"
"flag"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"github.com/gorilla/handlers"
"github.com/gorilla/mux"
"github.com/gorilla/websocket"
uuid "github.com/satori/go.uuid"
2019-01-07 11:45:13 +08:00
kafka "github.com/segmentio/kafka-go"
2019-01-05 16:06:47 +08:00
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" )
2019-01-07 11:45:13 +08:00
// consumer, err := kafka.NewConsumer(&kafka.ConfigMap{
// "bootstrap.servers": *bootstrapServers,
// "group.id": group,
// "session.timeout.ms": 6000,
// "auto.offset.reset": "latest"})
reader := kafka . NewReader ( kafka . ReaderConfig {
Brokers : [] string { * bootstrapServers },
Topic : topic ,
Partition : 0 ,
MinBytes : 10e1 ,
MaxBytes : 10e6 , // 10MB
})
reader . SetOffset ( - 1 ) //lastoffset
2019-01-05 16:06:47 +08:00
2019-01-07 11:45:13 +08:00
// 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)
// }
2019-01-05 16:06:47 +08:00
log . WithFields ( log . Fields {
2019-01-07 11:45:13 +08:00
"reader" : reader ,
2019-01-05 16:06:47 +08:00
}). Info ( "Created " )
2019-01-07 11:45:13 +08:00
// consumer.SubscribeTopics([]string{topic}, nil)
2019-01-05 16:06:47 +08:00
run := true
for run == true {
select {
case sig := <- sigchan :
log . WithFields ( log . Fields { "signal" : sig }). Info ( " terminating" )
run = false
log . Fatal ( "break " )
default :
2019-01-07 11:45:13 +08:00
// ev := consumer.Poll(100)
// if ev == nil {
// continue
// }
2019-01-05 16:06:47 +08:00
2019-01-07 11:45:13 +08:00
// 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")
// }
for {
m , err := reader . ReadMessage ( context . Background ())
if err != nil {
log . WithFields ( log . Fields { "error" : err }). Error ( "read error" )
2019-01-05 16:06:47 +08:00
}
2019-01-07 11:45:13 +08:00
// fmt.Printf("message at offset %d: %s = %s\n", m.Offset, string(m.Key), string(m.Value))
log . WithFields ( log . Fields { "offset" : m . Offset , "message" : string ( m . Value )}). Info ( "got message" )
c . hub . broadcast <- m . Value
2019-01-05 16:06:47 +08:00
}
}
}
2019-01-07 11:45:13 +08:00
// consumer.Close()
2019-01-05 16:06:47 +08:00
}
// 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 )
}