Live: include a streaming event manager (#26537)
This commit is contained in:
+5
-5
@@ -56,6 +56,7 @@ func (hs *HTTPServer) registerRoutes() {
|
||||
r.Get("/admin/orgs", reqGrafanaAdmin, hs.Index)
|
||||
r.Get("/admin/orgs/edit/:id", reqGrafanaAdmin, hs.Index)
|
||||
r.Get("/admin/stats", reqGrafanaAdmin, hs.Index)
|
||||
r.Get("/admin/live", reqGrafanaAdmin, hs.Index)
|
||||
r.Get("/admin/ldap", reqGrafanaAdmin, hs.Index)
|
||||
|
||||
r.Get("/styleguide", reqSignedIn, hs.Index)
|
||||
@@ -422,11 +423,10 @@ func (hs *HTTPServer) registerRoutes() {
|
||||
avatarCacheServer := avatar.NewCacheServer()
|
||||
r.Get("/avatar/:hash", avatarCacheServer.Handler)
|
||||
|
||||
// Websocket
|
||||
r.Any("/ws", hs.streamManager.Serve)
|
||||
|
||||
// streams
|
||||
//r.Post("/api/streams/push", reqSignedIn, bind(dtos.StreamMessage{}), liveConn.PushToStream)
|
||||
// Live streaming
|
||||
if hs.Live != nil {
|
||||
r.Any("/live/*", hs.Live.Handler)
|
||||
}
|
||||
|
||||
// Snapshots
|
||||
r.Post("/api/snapshots/", reqSnapshotPublicModeOrSignedIn, bind(models.CreateDashboardSnapshotCommand{}), CreateDashboardSnapshot)
|
||||
|
||||
+20
-9
@@ -10,11 +10,11 @@ import (
|
||||
"path"
|
||||
"sync"
|
||||
|
||||
"github.com/grafana/grafana/pkg/services/live"
|
||||
"github.com/grafana/grafana/pkg/services/search"
|
||||
|
||||
"github.com/grafana/grafana/pkg/plugins/backendplugin"
|
||||
|
||||
"github.com/grafana/grafana/pkg/api/live"
|
||||
"github.com/grafana/grafana/pkg/api/routing"
|
||||
httpstatic "github.com/grafana/grafana/pkg/api/static"
|
||||
"github.com/grafana/grafana/pkg/bus"
|
||||
@@ -48,12 +48,11 @@ func init() {
|
||||
}
|
||||
|
||||
type HTTPServer struct {
|
||||
log log.Logger
|
||||
macaron *macaron.Macaron
|
||||
context context.Context
|
||||
streamManager *live.StreamManager
|
||||
httpSrv *http.Server
|
||||
middlewares []macaron.Handler
|
||||
log log.Logger
|
||||
macaron *macaron.Macaron
|
||||
context context.Context
|
||||
httpSrv *http.Server
|
||||
middlewares []macaron.Handler
|
||||
|
||||
RouteRegister routing.RouteRegister `inject:""`
|
||||
Bus bus.Bus `inject:""`
|
||||
@@ -71,12 +70,25 @@ type HTTPServer struct {
|
||||
BackendPluginManager backendplugin.Manager `inject:""`
|
||||
PluginManager *plugins.PluginManager `inject:""`
|
||||
SearchService *search.SearchService `inject:""`
|
||||
Live *live.GrafanaLive
|
||||
}
|
||||
|
||||
func (hs *HTTPServer) Init() error {
|
||||
hs.log = log.New("http.server")
|
||||
|
||||
hs.streamManager = live.NewStreamManager()
|
||||
// Set up a websocket broker
|
||||
if hs.Cfg.IsLiveEnabled() { // feature flag
|
||||
node, err := live.InitalizeBroker()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
hs.Live = node
|
||||
|
||||
// Spit random walk to example
|
||||
go live.RunRandomCSV(hs.Live, "random-2s-stream", 2000, 0)
|
||||
go live.RunRandomCSV(hs.Live, "random-flakey-stream", 400, .6)
|
||||
}
|
||||
|
||||
hs.macaron = hs.newMacaron()
|
||||
hs.registerRoutes()
|
||||
|
||||
@@ -91,7 +103,6 @@ func (hs *HTTPServer) Run(ctx context.Context) error {
|
||||
hs.context = ctx
|
||||
|
||||
hs.applyRoutes()
|
||||
hs.streamManager.Run(ctx)
|
||||
|
||||
hs.httpSrv = &http.Server{
|
||||
Addr: fmt.Sprintf("%s:%s", setting.HttpAddr, setting.HttpPort),
|
||||
|
||||
@@ -342,6 +342,12 @@ func (hs *HTTPServer) setIndexViewData(c *models.ReqContext) (*dtos.IndexViewDat
|
||||
{Text: "Stats", Id: "server-stats", Url: setting.AppSubUrl + "/admin/stats", Icon: "graph-bar"},
|
||||
}
|
||||
|
||||
if hs.Live != nil {
|
||||
adminNavLinks = append(adminNavLinks, &dtos.NavLink{
|
||||
Text: "Live", Id: "live", Url: setting.AppSubUrl + "/admin/live", Icon: "water",
|
||||
})
|
||||
}
|
||||
|
||||
if setting.LDAPEnabled {
|
||||
adminNavLinks = append(adminNavLinks, &dtos.NavLink{
|
||||
Text: "LDAP", Id: "ldap", Url: setting.AppSubUrl + "/admin/ldap", Icon: "book",
|
||||
|
||||
@@ -1,130 +0,0 @@
|
||||
package live
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/gorilla/websocket"
|
||||
"github.com/grafana/grafana/pkg/components/simplejson"
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
)
|
||||
|
||||
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 upgrader = websocket.Upgrader{
|
||||
ReadBufferSize: 1024,
|
||||
WriteBufferSize: 1024,
|
||||
CheckOrigin: func(r *http.Request) bool {
|
||||
return true
|
||||
},
|
||||
}
|
||||
|
||||
type connection struct {
|
||||
hub *hub
|
||||
ws *websocket.Conn
|
||||
send chan []byte
|
||||
log log.Logger
|
||||
}
|
||||
|
||||
func newConnection(ws *websocket.Conn, hub *hub, logger log.Logger) *connection {
|
||||
return &connection{
|
||||
hub: hub,
|
||||
send: make(chan []byte, 256),
|
||||
ws: ws,
|
||||
log: logger,
|
||||
}
|
||||
}
|
||||
|
||||
func (c *connection) readPump() {
|
||||
defer func() {
|
||||
c.hub.unregister <- c
|
||||
c.ws.Close()
|
||||
}()
|
||||
|
||||
c.ws.SetReadLimit(maxMessageSize)
|
||||
if err := c.ws.SetReadDeadline(time.Now().Add(pongWait)); err != nil {
|
||||
c.log.Warn("Setting read deadline failed", "err", err)
|
||||
}
|
||||
c.ws.SetPongHandler(func(string) error {
|
||||
return c.ws.SetReadDeadline(time.Now().Add(pongWait))
|
||||
})
|
||||
for {
|
||||
_, message, err := c.ws.ReadMessage()
|
||||
if err != nil {
|
||||
if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway) {
|
||||
c.log.Info("error", "err", err)
|
||||
}
|
||||
break
|
||||
}
|
||||
|
||||
c.handleMessage(message)
|
||||
}
|
||||
}
|
||||
|
||||
func (c *connection) handleMessage(message []byte) {
|
||||
json, err := simplejson.NewJson(message)
|
||||
if err != nil {
|
||||
log.Errorf(3, "Unreadable message on websocket channel. error: %v", err)
|
||||
}
|
||||
|
||||
msgType := json.Get("action").MustString()
|
||||
streamName := json.Get("stream").MustString()
|
||||
|
||||
if len(streamName) == 0 {
|
||||
log.Errorf(3, "Not allowed to subscribe to empty stream name")
|
||||
return
|
||||
}
|
||||
|
||||
switch msgType {
|
||||
case "subscribe":
|
||||
c.hub.subChannel <- &streamSubscription{name: streamName, conn: c}
|
||||
case "unsubscribe":
|
||||
c.hub.subChannel <- &streamSubscription{name: streamName, conn: c, remove: true}
|
||||
}
|
||||
}
|
||||
|
||||
func (c *connection) write(mt int, payload []byte) error {
|
||||
if err := c.ws.SetWriteDeadline(time.Now().Add(writeWait)); err != nil {
|
||||
return err
|
||||
}
|
||||
return c.ws.WriteMessage(mt, payload)
|
||||
}
|
||||
|
||||
// writePump pumps messages from the hub to the websocket connection.
|
||||
func (c *connection) writePump() {
|
||||
ticker := time.NewTicker(pingPeriod)
|
||||
defer func() {
|
||||
ticker.Stop()
|
||||
c.ws.Close()
|
||||
}()
|
||||
for {
|
||||
select {
|
||||
case message, ok := <-c.send:
|
||||
if !ok {
|
||||
if err := c.write(websocket.CloseMessage, []byte{}); err != nil {
|
||||
c.log.Warn("Failed to write close message to connection", "err", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
if err := c.write(websocket.TextMessage, message); err != nil {
|
||||
return
|
||||
}
|
||||
case <-ticker.C:
|
||||
if err := c.write(websocket.PingMessage, []byte{}); err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,99 +0,0 @@
|
||||
package live
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/grafana/grafana/pkg/api/dtos"
|
||||
"github.com/grafana/grafana/pkg/components/simplejson"
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
)
|
||||
|
||||
type hub struct {
|
||||
log log.Logger
|
||||
connections map[*connection]bool
|
||||
streams map[string]map[*connection]bool
|
||||
|
||||
register chan *connection
|
||||
unregister chan *connection
|
||||
streamChannel chan *dtos.StreamMessage
|
||||
subChannel chan *streamSubscription
|
||||
}
|
||||
|
||||
type streamSubscription struct {
|
||||
conn *connection
|
||||
name string
|
||||
remove bool
|
||||
}
|
||||
|
||||
func newHub() *hub {
|
||||
return &hub{
|
||||
connections: make(map[*connection]bool),
|
||||
streams: make(map[string]map[*connection]bool),
|
||||
register: make(chan *connection),
|
||||
unregister: make(chan *connection),
|
||||
streamChannel: make(chan *dtos.StreamMessage),
|
||||
subChannel: make(chan *streamSubscription),
|
||||
log: log.New("stream.hub"),
|
||||
}
|
||||
}
|
||||
|
||||
func (h *hub) run(ctx context.Context) {
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
case c := <-h.register:
|
||||
h.connections[c] = true
|
||||
h.log.Info("New connection", "total", len(h.connections))
|
||||
|
||||
case c := <-h.unregister:
|
||||
if _, ok := h.connections[c]; ok {
|
||||
h.log.Info("Closing connection", "total", len(h.connections))
|
||||
delete(h.connections, c)
|
||||
close(c.send)
|
||||
}
|
||||
// hand stream subscriptions
|
||||
case sub := <-h.subChannel:
|
||||
h.log.Info("Subscribing", "channel", sub.name, "remove", sub.remove)
|
||||
subscribers, exists := h.streams[sub.name]
|
||||
|
||||
// handle unsubscribe
|
||||
if exists && sub.remove {
|
||||
delete(subscribers, sub.conn)
|
||||
continue
|
||||
}
|
||||
|
||||
if !exists {
|
||||
subscribers = make(map[*connection]bool)
|
||||
h.streams[sub.name] = subscribers
|
||||
}
|
||||
|
||||
subscribers[sub.conn] = true
|
||||
|
||||
// handle stream messages
|
||||
case message := <-h.streamChannel:
|
||||
subscribers, exists := h.streams[message.Stream]
|
||||
if !exists || len(subscribers) == 0 {
|
||||
h.log.Info("Message to stream without subscribers", "stream", message.Stream)
|
||||
continue
|
||||
}
|
||||
|
||||
messageBytes, _ := simplejson.NewFromAny(message).Encode()
|
||||
for sub := range subscribers {
|
||||
// check if channel is open
|
||||
if _, ok := h.connections[sub]; !ok {
|
||||
delete(subscribers, sub)
|
||||
continue
|
||||
}
|
||||
|
||||
select {
|
||||
case sub.send <- messageBytes:
|
||||
default:
|
||||
close(sub.send)
|
||||
delete(h.connections, sub)
|
||||
delete(subscribers, sub)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,102 +0,0 @@
|
||||
package live
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net/http"
|
||||
"sync"
|
||||
|
||||
"github.com/grafana/grafana/pkg/components/simplejson"
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
"github.com/grafana/grafana/pkg/models"
|
||||
)
|
||||
|
||||
type StreamManager struct {
|
||||
log log.Logger
|
||||
streams map[string]*Stream
|
||||
streamRWMutex *sync.RWMutex
|
||||
hub *hub
|
||||
}
|
||||
|
||||
func NewStreamManager() *StreamManager {
|
||||
return &StreamManager{
|
||||
hub: newHub(),
|
||||
log: log.New("stream.manager"),
|
||||
streams: make(map[string]*Stream),
|
||||
streamRWMutex: &sync.RWMutex{},
|
||||
}
|
||||
}
|
||||
|
||||
func (sm *StreamManager) Run(context context.Context) {
|
||||
log.Debugf("Initializing Stream Manager")
|
||||
|
||||
go func() {
|
||||
sm.hub.run(context)
|
||||
log.Infof("Stopped Stream Manager")
|
||||
}()
|
||||
}
|
||||
|
||||
func (sm *StreamManager) Serve(w http.ResponseWriter, r *http.Request) {
|
||||
sm.log.Info("Upgrading to WebSocket")
|
||||
|
||||
ws, err := upgrader.Upgrade(w, r, nil)
|
||||
if err != nil {
|
||||
sm.log.Error("Failed to upgrade connection to WebSocket", "error", err)
|
||||
return
|
||||
}
|
||||
|
||||
c := newConnection(ws, sm.hub, sm.log)
|
||||
sm.hub.register <- c
|
||||
|
||||
go c.writePump()
|
||||
c.readPump()
|
||||
}
|
||||
|
||||
func (s *StreamManager) GetStreamList() models.StreamList {
|
||||
list := make(models.StreamList, 0)
|
||||
|
||||
for _, stream := range s.streams {
|
||||
list = append(list, &models.StreamInfo{
|
||||
Name: stream.name,
|
||||
})
|
||||
}
|
||||
|
||||
return list
|
||||
}
|
||||
|
||||
func (s *StreamManager) Push(packet *models.StreamPacket) {
|
||||
stream, exist := s.streams[packet.Stream]
|
||||
|
||||
if !exist {
|
||||
s.log.Info("Creating metric stream", "name", packet.Stream)
|
||||
stream = NewStream(packet.Stream)
|
||||
s.streams[stream.name] = stream
|
||||
}
|
||||
|
||||
stream.Push(packet)
|
||||
}
|
||||
|
||||
type Stream struct {
|
||||
subscribers []*connection
|
||||
name string
|
||||
}
|
||||
|
||||
func NewStream(name string) *Stream {
|
||||
return &Stream{
|
||||
subscribers: make([]*connection, 0),
|
||||
name: name,
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Stream) Push(packet *models.StreamPacket) {
|
||||
messageBytes, _ := simplejson.NewFromAny(packet).Encode()
|
||||
|
||||
for _, sub := range s.subscribers {
|
||||
// check if channel is open
|
||||
// if _, ok := h.connections[sub]; !ok {
|
||||
// delete(s.subscribers, sub)
|
||||
// continue
|
||||
// }
|
||||
|
||||
sub.send <- messageBytes
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user