121 lines
2.7 KiB
Go
121 lines
2.7 KiB
Go
package activity
|
|
|
|
import (
|
|
"crypto/rand"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"fmt"
|
|
"net/http"
|
|
"strings"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"github.com/local/glpi-neural-brain/internal/model"
|
|
)
|
|
|
|
type Broker struct {
|
|
mu sync.RWMutex
|
|
next int
|
|
subs map[int]chan model.Activity
|
|
recent []model.Activity
|
|
maxRecent int
|
|
sinkMu sync.RWMutex
|
|
sink func(model.Activity)
|
|
eventSeq atomic.Uint64
|
|
eventNonce string
|
|
}
|
|
|
|
func New(maxRecent int) *Broker {
|
|
if maxRecent < 10 {
|
|
maxRecent = 100
|
|
}
|
|
var raw [6]byte
|
|
_, _ = rand.Read(raw[:])
|
|
nonce := hex.EncodeToString(raw[:])
|
|
if strings.Trim(nonce, "0") == "" {
|
|
nonce = fmt.Sprintf("%x", time.Now().UnixNano())
|
|
}
|
|
return &Broker{subs: map[int]chan model.Activity{}, maxRecent: maxRecent, eventNonce: nonce}
|
|
}
|
|
|
|
// SetSink attaches a non-UI activity recorder. The sink is called after the
|
|
// in-memory broker lock is released, so recording cannot block SSE delivery or
|
|
// deadlock the broker.
|
|
func (b *Broker) SetSink(sink func(model.Activity)) {
|
|
b.sinkMu.Lock()
|
|
b.sink = sink
|
|
b.sinkMu.Unlock()
|
|
}
|
|
|
|
func (b *Broker) Publish(a model.Activity) {
|
|
if a.Timestamp.IsZero() {
|
|
a.Timestamp = time.Now().UTC()
|
|
}
|
|
if a.ID == "" {
|
|
a.ID = fmt.Sprintf("evt-%s-%d-%d", b.eventNonce, a.Timestamp.UnixNano(), b.eventSeq.Add(1))
|
|
}
|
|
b.mu.Lock()
|
|
b.recent = append(b.recent, a)
|
|
if len(b.recent) > b.maxRecent {
|
|
b.recent = append([]model.Activity(nil), b.recent[len(b.recent)-b.maxRecent:]...)
|
|
}
|
|
for _, ch := range b.subs {
|
|
select {
|
|
case ch <- a:
|
|
default:
|
|
}
|
|
}
|
|
b.mu.Unlock()
|
|
b.sinkMu.RLock()
|
|
sink := b.sink
|
|
b.sinkMu.RUnlock()
|
|
if sink != nil {
|
|
sink(a)
|
|
}
|
|
}
|
|
func (b *Broker) Recent() []model.Activity {
|
|
b.mu.RLock()
|
|
defer b.mu.RUnlock()
|
|
return append([]model.Activity(nil), b.recent...)
|
|
}
|
|
func (b *Broker) ServeSSE(w http.ResponseWriter, r *http.Request) {
|
|
fl, ok := w.(http.Flusher)
|
|
if !ok {
|
|
http.Error(w, "streaming unsupported", 500)
|
|
return
|
|
}
|
|
w.Header().Set("Content-Type", "text/event-stream")
|
|
w.Header().Set("Cache-Control", "no-cache")
|
|
w.Header().Set("Connection", "keep-alive")
|
|
b.mu.Lock()
|
|
id := b.next
|
|
b.next++
|
|
ch := make(chan model.Activity, 64)
|
|
b.subs[id] = ch
|
|
recent := append([]model.Activity(nil), b.recent...)
|
|
b.mu.Unlock()
|
|
defer func() { b.mu.Lock(); delete(b.subs, id); close(ch); b.mu.Unlock() }()
|
|
enc := func(a model.Activity) {
|
|
data, _ := json.Marshal(a)
|
|
fmt.Fprintf(w, "event: activity\ndata: %s\n\n", data)
|
|
fl.Flush()
|
|
}
|
|
for _, a := range recent {
|
|
enc(a)
|
|
}
|
|
ticker := time.NewTicker(20 * time.Second)
|
|
defer ticker.Stop()
|
|
for {
|
|
select {
|
|
case <-r.Context().Done():
|
|
return
|
|
case a := <-ch:
|
|
enc(a)
|
|
case <-ticker.C:
|
|
fmt.Fprint(w, ": ping\n\n")
|
|
fl.Flush()
|
|
}
|
|
}
|
|
}
|