2019-01-05 16:06:47 +08:00
package main
import (
"context"
"flag"
"net/http"
"os"
"os/signal"
2019-01-07 14:29:04 +08:00
"strings"
2019-01-05 16:06:47 +08:00
"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' }
addr = flag . String ( "addr" , ":9000" , "http service address" )
bootstrapServers = flag . String ( "bootstrap-servers" , "localhost:9092" , "kafka bootstrap servers" )
2019-01-07 11:54:11 +08:00
partition = flag . Int ( "partition" , 0 , "kafka tipic partion" )
wait = flag . Duration ( "graceful-timeout" , time . Second * 15 , "the duration for which the server gracefully wait for existing connections to finish - e.g. 15s or 1m" )
2019-01-05 16:06:47 +08:00
upgrader = websocket . Upgrader {
ReadBufferSize : 1024 ,
WriteBufferSize : 1024 ,
}
2019-01-08 16:14:11 +08:00
hubMap = make ( map [ string ] Hub )
run = true
2019-01-05 16:06:47 +08:00
)
// 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
}
2019-01-08 16:14:11 +08:00
func create ( topic string ) * kafka . Reader {
group , _ := uuid . NewV4 ()
2019-01-08 10:50:41 +08:00
2019-01-05 16:06:47 +08:00
log . WithFields ( log . Fields {
"topic" : topic ,
"bootstrap servers" : * bootstrapServers ,
2019-01-08 10:50:41 +08:00
"partition" : * partition ,
"group" : group }). Info ( "creating kafka client ... " )
2019-01-05 16:06:47 +08:00
2019-01-07 11:45:13 +08:00
reader := kafka . NewReader ( kafka . ReaderConfig {
2019-01-07 14:29:04 +08:00
Brokers : strings . Split ( * bootstrapServers , "," ),
2019-01-07 11:45:13 +08:00
Topic : topic ,
2019-01-07 14:29:04 +08:00
Partition : * partition ,
2019-01-07 11:45:13 +08:00
MinBytes : 10e1 ,
MaxBytes : 10e6 , // 10MB
})
2019-01-07 14:29:04 +08:00
reader . SetOffset ( - 1 ) //latest offset
2019-01-05 16:06:47 +08:00
2019-01-08 16:14:11 +08:00
return reader
}
2019-01-05 16:06:47 +08:00
2019-01-08 16:14:11 +08:00
// 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 () {
defer func () {
c . hub . unregister <- c
c . conn . Close ()
log . WithField ( "client" , c ). Warn ( "exit" )
}()
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 )
2019-01-05 16:06:47 +08:00
for run == true {
select {
case sig := <- sigchan :
log . WithFields ( log . Fields { "signal" : sig }). Info ( " terminating" )
run = false
log . Fatal ( "break " )
2019-01-08 16:14:11 +08:00
// default:
// for {
// m, err := c.hub.reader.ReadMessage(context.Background())
// if err != nil {
// log.WithFields(log.Fields{"error": err}).Error("read error")
// }
// 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
}
}
}
// 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.
2019-01-08 16:14:11 +08:00
func serveWs ( w http . ResponseWriter , r * http . Request ) {
2019-01-05 16:06:47 +08:00
vars := mux . Vars ( r )
topic := vars [ "topic" ]
conn , err := upgrader . Upgrade ( w , r , nil )
if err != nil {
log . Error ( "upgrade connection failed " , err )
return
}
2019-01-08 16:14:11 +08:00
h , ok := hubMap [ topic ]
if ! ok {
log . WithField ( "topic" , topic ). Info ( "create new hub " )
h = * newHub ( topic )
go h . run ()
hubMap [ topic ] = h
} else {
log . Info ( "join hub" )
}
client := & Client { hub : & h , conn : conn , send : make ( chan [] byte , 256 )}
2019-01-05 16:06:47 +08:00
client . hub . register <- client
// Allow collection of memory referenced by the caller by doing all work in
// new goroutines.
go client . writePump ()
2019-01-08 16:14:11 +08:00
go client . readPump ()
2019-01-05 16:06:47 +08:00
}
// Hub maintains the set of active clients and broadcasts messages to the
// clients.
type Hub struct {
2019-01-08 16:14:11 +08:00
reader * kafka . Reader
2019-01-05 16:06:47 +08:00
// 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
}
2019-01-08 16:14:11 +08:00
func newHub ( topic string ) * Hub {
2019-01-05 16:06:47 +08:00
return & Hub {
2019-01-08 16:14:11 +08:00
reader : create ( topic ),
2019-01-05 16:06:47 +08:00
broadcast : make ( chan [] byte ),
register : make ( chan * Client ),
unregister : make ( chan * Client ),
clients : make ( map [ * Client ] bool ),
}
}
2019-01-08 16:14:11 +08:00
func ( h * Hub ) read () {
log . WithField ( "reader" , h . reader ). Info ( "start to read" )
for run == true {
m , err := h . reader . ReadMessage ( context . Background ())
if err != nil {
log . WithFields ( log . Fields { "error" : err }). Error ( "read error" )
}
log . WithFields ( log . Fields { "offset" : m . Offset , "message" : string ( m . Value )}). Info ( "got message" )
h . broadcast <- m . Value
}
}
2019-01-05 16:06:47 +08:00
func ( h * Hub ) run () {
2019-01-08 16:14:11 +08:00
go h . read ()
for run == true {
2019-01-05 16:06:47 +08:00
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 () {
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.
}
2019-01-08 16:14:11 +08:00
// hub := newHub()
// go hub.run()
2019-01-05 16:06:47 +08:00
r . HandleFunc ( "/ws/{topic}" , func ( w http . ResponseWriter , r * http . Request ) {
2019-01-08 16:14:11 +08:00
serveWs ( w , r )
2019-01-05 16:06:47 +08:00
})
// 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 {
2019-01-08 16:14:11 +08:00
log . WithFields ( log . Fields { "error" : err }). Error ( "error" )
2019-01-05 16:06:47 +08:00
}
}()
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.
2019-01-07 11:54:11 +08:00
ctx , cancel := context . WithTimeout ( context . Background (), * wait )
2019-01-05 16:06:47 +08:00
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 )
}