Files

707 lines
18 KiB
Go

package main
import (
"encoding/json"
"fmt"
"io/fs"
"log"
"net/http"
"os"
"sync"
"time"
"exchange-monitor/db"
)
// ============================================================
// SSE Hub — manages connected browser clients
// ============================================================
type sseClient struct {
ch chan []byte
done chan struct{}
filter string
}
type SSEHub struct {
mu sync.RWMutex
clients map[*sseClient]bool
seq uint64
}
func NewSSEHub() *SSEHub {
return &SSEHub{
clients: make(map[*sseClient]bool),
}
}
func (h *SSEHub) Subscribe(filter string) *sseClient {
c := &sseClient{
ch: make(chan []byte, 64),
done: make(chan struct{}),
filter: filter,
}
h.mu.Lock()
h.clients[c] = true
h.mu.Unlock()
log.Printf("[Web] SSE client connected (clients=%d)", len(h.clients))
return c
}
func (h *SSEHub) Unsubscribe(c *sseClient) {
h.mu.Lock()
delete(h.clients, c)
count := len(h.clients)
h.mu.Unlock()
close(c.done)
log.Printf("[Web] SSE client disconnected (clients=%d)", count)
}
func (h *SSEHub) Broadcast(event string, data interface{}) {
raw, err := json.Marshal(map[string]interface{}{
"event": event,
"data": data,
"ts": time.Now().UnixMilli(),
})
if err != nil {
return
}
h.mu.RLock()
defer h.mu.RUnlock()
for c := range h.clients {
select {
case c.ch <- raw:
default:
}
}
}
// ============================================================
// History Ring Buffers
// ============================================================
const maxHistoryPoints = 500
type pricePoint struct {
T int64 `json:"t"`
P float64 `json:"p"`
}
type priceHistory struct {
mu sync.RWMutex
buffers map[string]map[string][]pricePoint
}
func newPriceHistory() *priceHistory {
return &priceHistory{
buffers: make(map[string]map[string][]pricePoint),
}
}
func (ph *priceHistory) Record(coin, exchange string, price float64) {
ph.mu.Lock()
defer ph.mu.Unlock()
if ph.buffers[coin] == nil {
ph.buffers[coin] = make(map[string][]pricePoint)
}
buf := ph.buffers[coin][exchange]
buf = append(buf, pricePoint{T: time.Now().UnixMilli(), P: price})
if len(buf) > maxHistoryPoints {
buf = buf[len(buf)-maxHistoryPoints:]
}
ph.buffers[coin][exchange] = buf
}
func (ph *priceHistory) GetHistory(coin, exchange string, limit int) []pricePoint {
ph.mu.RLock()
defer ph.mu.RUnlock()
buf := ph.buffers[coin][exchange]
if len(buf) == 0 {
return nil
}
if limit <= 0 || limit >= len(buf) {
r := make([]pricePoint, len(buf))
copy(r, buf)
return r
}
r := make([]pricePoint, limit)
copy(r, buf[len(buf)-limit:])
return r
}
// ============================================================
// Spread History — tracks 3-exchange max spread % per coin
// ============================================================
type spreadPoint struct {
T int64 `json:"t"`
Spread float64 `json:"s"` // 3-exchange max spread %
}
type spreadHistory struct {
mu sync.RWMutex
buffers map[string][]spreadPoint
}
func newSpreadHistory() *spreadHistory {
return &spreadHistory{
buffers: make(map[string][]spreadPoint),
}
}
func (sh *spreadHistory) Record(coin string, spread float64) {
sh.mu.Lock()
defer sh.mu.Unlock()
sh.buffers[coin] = append(sh.buffers[coin], spreadPoint{T: time.Now().UnixMilli(), Spread: spread})
if len(sh.buffers[coin]) > maxHistoryPoints {
sh.buffers[coin] = sh.buffers[coin][len(sh.buffers[coin])-maxHistoryPoints:]
}
}
func (sh *spreadHistory) GetHistory(coin string, limit int) []spreadPoint {
sh.mu.RLock()
defer sh.mu.RUnlock()
buf := sh.buffers[coin]
if len(buf) == 0 {
return nil
}
if limit <= 0 || limit >= len(buf) {
r := make([]spreadPoint, len(buf))
copy(r, buf)
return r
}
r := make([]spreadPoint, limit)
copy(r, buf[len(buf)-limit:])
return r
}
// ============================================================
// Dashboard
// ============================================================
type Dashboard struct {
hub *SSEHub
history *priceHistory
spreads *spreadHistory
store *PriceStore
db *db.DB
addr string
cfg *Config
// cached scan results
mu sync.RWMutex
lastScan []ThreeExSpread
scanTime time.Time
// connection status — exchange -> last update time
connMu sync.RWMutex
connMap map[string]time.Time
// Momentum tracker
momentumTracker *MomentumTracker
// Trend detector
trendDetector *TrendDetector
// Cumulative tracker (1m/5m multi-exchange consensus)
cumulativeTracker *CumulativeTracker
// Trend filter (K-line based quiet + EMA filter)
trendFilter *TrendFilter
// Surge detector
surgeDetector *SurgeDetector
// Binance momentum detector
binanceMomentumDetector *BinanceMomentumDetector
}
func NewDashboard(store *PriceStore, database *db.DB, addr string, cfg *Config, momentumTracker *MomentumTracker, trendDetector *TrendDetector, cumulativeTracker *CumulativeTracker, trendFilter *TrendFilter, surgeDetector *SurgeDetector, binanceMomentumDetector *BinanceMomentumDetector) *Dashboard {
d := &Dashboard{
hub: NewSSEHub(),
history: newPriceHistory(),
spreads: newSpreadHistory(),
store: store,
db: database,
addr: addr,
cfg: cfg,
connMap: make(map[string]time.Time),
momentumTracker: momentumTracker,
trendDetector: trendDetector,
cumulativeTracker: cumulativeTracker,
trendFilter: trendFilter,
surgeDetector: surgeDetector,
binanceMomentumDetector: binanceMomentumDetector,
}
// Wire trend event persistence to SQLite
if trendDetector != nil && database != nil {
trendDetector.OnEvent = func(ev TrendEvent) {
database.InsertTrendEvent(ev.Coin, ev.PrevState, ev.NewState, ev.Direction,
ev.ZScore, ev.Volatility, ev.BGChange, ev.HLChange, ev.BNChange, ev.OKXChange,
ev.ExAgree, ev.ExTotal)
}
}
// Wire trend filter signal broadcast via SSE
if trendFilter != nil {
trendFilter.OnNewSignal = func(sig TrendSignal) {
d.hub.Broadcast("trend_signal", sig)
}
}
// Wire cumulative event persistence to SQLite
if cumulativeTracker != nil && database != nil {
cumulativeTracker.OnEvent = func(ev CmEvent) {
database.InsertCmEvent(ev.Coin, ev.PrevState, ev.NewState, ev.Direction,
ev.Score, ev.AvgChange, ev.ExAgree, ev.ExTotal,
ev.BGChange1m, ev.HLChange1m, ev.BNChange1m, ev.OKXChange1m,
ev.BGChange5m, ev.HLChange5m, ev.BNChange5m, ev.OKXChange5m)
}
}
return d
}
func (d *Dashboard) Run() {
go d.broadcastLoop()
mux := http.NewServeMux()
// Try disk-based serving first (hot-reload friendly), fall back to embed
staticSub, err := fs.Sub(staticFS, "frontend/dist")
if diskFS := os.DirFS("frontend/dist"); true {
if _, diskErr := fs.Stat(diskFS, "index.html"); diskErr == nil {
staticSub = diskFS
}
}
if err == nil && staticSub != nil {
mux.Handle("GET /static/", http.StripPrefix("/static/", http.FileServer(http.FS(staticSub))))
}
mux.HandleFunc("GET /", d.handleIndex)
mux.HandleFunc("GET /api/status", d.handleStatus)
mux.HandleFunc("GET /api/history", d.handleHistory)
mux.HandleFunc("GET /api/spread-history", d.handleSpreadHistory)
mux.HandleFunc("GET /api/connections", d.handleConnStatus)
mux.HandleFunc("GET /api/trend-history", d.handleTrendHistory)
mux.HandleFunc("GET /api/cm-history", d.handleCmHistory)
mux.HandleFunc("GET /api/trend-signals", d.handleTrendSignals)
mux.HandleFunc("GET /api/surge-events", d.handleSurgeEvents)
mux.HandleFunc("GET /events", d.handleSSE)
mux.HandleFunc("GET /binance", d.handleBinanceIndex)
mux.HandleFunc("GET /api/binance-alerts", d.handleBinanceAlerts)
server := &http.Server{
Addr: d.addr,
Handler: mux,
ReadTimeout: 10 * time.Second,
WriteTimeout: 0,
}
log.Printf("[Web] Dashboard listening on http://%s", d.addr)
if err := server.ListenAndServe(); err != nil && err != http.ErrServerClosed {
log.Printf("[Web] Server error: %v", err)
}
}
// broadcastLoop pushes data to SSE clients every 1 second.
func (d *Dashboard) broadcastLoop() {
tick := time.NewTicker(1 * time.Second)
defer tick.Stop()
for range tick.C {
snap := d.store.GetAll()
if len(snap) == 0 {
continue
}
// 1. Prices + 3-exchange spreads
var prices []map[string]interface{}
for _, coin := range TrackedCoins {
exMap := snap[coin.Name]
if exMap == nil {
continue
}
entry := map[string]interface{}{
"coin": coin.Name,
}
for ex, p := range exMap {
entry[ex] = p
}
for ex := range exMap {
sp := d.store.GetSpread(coin.Name, ex)
if sp > 0 {
entry[ex+"_spread"] = sp
}
}
// Calculate 3-exchange max spread
bnP := exMap[ExBinance]
okxP := exMap[ExOKX]
bgP := exMap[ExBitget]
if bnP > 0 && okxP > 0 && bgP > 0 {
prices_ := []float64{bnP, okxP, bgP}
minP, maxP := prices_[0], prices_[0]
for _, p := range prices_[1:] {
if p < minP { minP = p }
if p > maxP { maxP = p }
}
spreadPct := (maxP - minP) / minP * 100
entry["spread_3ex"] = spreadPct
d.spreads.Record(coin.Name, spreadPct)
}
prices = append(prices, entry)
}
d.hub.Broadcast("prices", prices)
// 2. 3-exchange scan results
d.mu.RLock()
scanCopy := d.lastScan
d.mu.RUnlock()
if len(scanCopy) > 0 {
scanList := make([]map[string]interface{}, 0, len(scanCopy))
for _, s := range scanCopy {
scanList = append(scanList, map[string]interface{}{
"coin": s.Coin,
"spread_pct": s.SpreadPct,
"bn_price": s.BnPrice,
"okx_price": s.OkxPrice,
"bg_price": s.BgPrice,
"max_ex": s.MaxEx,
"min_ex": s.MinEx,
})
}
d.hub.Broadcast("spread_3ex", scanList)
}
// 3. Connection status
d.connMu.RLock()
connInfo := make(map[string]string)
for ex, lastTime := range d.connMap {
age := time.Since(lastTime)
if age < 10*time.Second {
connInfo[ex] = "online"
} else if age < 30*time.Second {
connInfo[ex] = "stale"
} else {
connInfo[ex] = "offline"
}
}
d.connMu.RUnlock()
status := map[string]interface{}{
"coins": len(prices),
"connections": connInfo,
}
d.hub.Broadcast("status", status)
// 4. Momentum data (if enabled)
if d.momentumTracker != nil && d.cfg.MomentumEnabled {
momentumData := d.momentumTracker.Snapshot(d.cfg.MomentumThresholdPct)
if len(momentumData) > 0 {
d.hub.Broadcast("momentum", momentumData)
}
}
// 5. Trend detection (if enabled)
if d.trendDetector != nil && d.cfg.TrendEnabled {
d.trendDetector.Tick()
trendData := d.trendDetector.Snapshot()
if len(trendData) > 0 {
d.hub.Broadcast("trend", trendData)
}
}
// 6. Cumulative change tracking
if d.cumulativeTracker != nil {
d.cumulativeTracker.Tick()
cmData := d.cumulativeTracker.GetTopCoins(30)
if len(cmData) > 0 {
d.hub.Broadcast("cumulative", cmData)
}
}
// 7. Trend filter (K-line based quiet + EMA)
if d.trendFilter != nil {
d.trendFilter.Tick()
filterData := d.trendFilter.Snapshot(0)
if len(filterData) > 0 {
d.hub.Broadcast("trend_filter", filterData)
}
}
// 8. Surge status (current spread/baseline for all coins)
if d.surgeDetector != nil && d.cfg.SurgeEnabled {
surgeSnap := d.surgeDetector.Snapshot()
if len(surgeSnap) > 0 {
d.hub.Broadcast("surge", surgeSnap)
}
}
// 9. Binance momentum snapshot (1-minute change for all coins)
if d.binanceMomentumDetector != nil && d.cfg.BinanceMomentumEnabled {
bmSnap := d.binanceMomentumDetector.Snapshot()
if len(bmSnap) > 0 {
d.hub.Broadcast("binance_momentum", bmSnap)
}
}
}
}
// ============================================================
// Public methods called from main.go
// ============================================================
func (d *Dashboard) UpdateScan(spreads []ThreeExSpread) {
d.mu.Lock()
d.lastScan = spreads
d.scanTime = time.Now()
d.mu.Unlock()
}
func (d *Dashboard) RecordPrice(coin, exchange string, price float64) {
d.history.Record(coin, exchange, price)
}
// RecordConnStatus updates the last-seen time for an exchange.
func (d *Dashboard) RecordConnStatus(exchange string) {
d.connMu.Lock()
d.connMap[exchange] = time.Now()
d.connMu.Unlock()
}
// BroadcastEvent sends an immediate SSE event.
func (d *Dashboard) BroadcastEvent(event string, data interface{}) {
d.hub.Broadcast(event, data)
}
// ============================================================
// HTTP Handlers
// ============================================================
func (d *Dashboard) handleIndex(w http.ResponseWriter, r *http.Request) {
var data []byte
var err error
// Try disk first (hot reload)
data, err = os.ReadFile("frontend/dist/index.html")
if err != nil {
data, err = staticFS.ReadFile("frontend/dist/index.html")
}
if err != nil {
http.Error(w, "Not found", 404)
return
}
w.Header().Set("Content-Type", "text/html; charset=utf-8")
w.Write(data)
}
func (d *Dashboard) handleStatus(w http.ResponseWriter, r *http.Request) {
snap := d.store.GetAll()
resp := map[string]interface{}{
"prices": snap,
"coins": len(snap),
}
writeJSON(w, resp)
}
func (d *Dashboard) handleHistory(w http.ResponseWriter, r *http.Request) {
coin := r.URL.Query().Get("coin")
exchange := r.URL.Query().Get("exchange")
if coin == "" || exchange == "" {
snap := d.store.GetAll()
coins := make([]string, 0, len(snap))
for c := range snap {
coins = append(coins, c)
}
writeJSON(w, map[string]interface{}{"coins": coins, "exchanges": []string{ExBinance, ExOKX, ExBitget}})
return
}
points := d.history.GetHistory(coin, exchange, 300)
writeJSON(w, map[string]interface{}{
"coin": coin,
"exchange": exchange,
"points": points,
})
}
// handleSpreadHistory returns 3-exchange max spread history for a coin.
func (d *Dashboard) handleSpreadHistory(w http.ResponseWriter, r *http.Request) {
coin := r.URL.Query().Get("coin")
if coin == "" {
writeJSON(w, map[string]interface{}{"coins": trackedCoinNames()})
return
}
points := d.spreads.GetHistory(coin, 300)
writeJSON(w, map[string]interface{}{
"coin": coin,
"points": points,
})
}
// handleConnStatus returns connection health for all exchanges.
func (d *Dashboard) handleConnStatus(w http.ResponseWriter, r *http.Request) {
d.connMu.RLock()
conns := make(map[string]string)
for ex, t := range d.connMap {
age := time.Since(t)
switch {
case age < 10*time.Second:
conns[ex] = "online"
case age < 30*time.Second:
conns[ex] = "stale"
default:
conns[ex] = "offline"
}
}
d.connMu.RUnlock()
writeJSON(w, conns)
}
func (d *Dashboard) handleTrendHistory(w http.ResponseWriter, r *http.Request) {
var events interface{}
if d.db != nil {
records, err := d.db.GetTrendEvents(200)
if err == nil {
events = records
}
}
if events == nil {
if d.trendDetector != nil {
events = d.trendDetector.GetEvents(200)
} else {
events = []interface{}{}
}
}
writeJSON(w, map[string]interface{}{"events": events})
}
func (d *Dashboard) handleCmHistory(w http.ResponseWriter, r *http.Request) {
var events interface{}
if d.db != nil {
records, err := d.db.GetCmEvents(200)
if err == nil {
events = records
}
}
if events == nil {
if d.cumulativeTracker != nil {
events = d.cumulativeTracker.GetEvents(200)
} else {
events = []interface{}{}
}
}
writeJSON(w, map[string]interface{}{"events": events})
}
func (d *Dashboard) handleTrendSignals(w http.ResponseWriter, r *http.Request) {
var signals []TrendSignal
if d.trendFilter != nil {
signals = d.trendFilter.GetSignals(100)
}
if signals == nil {
signals = []TrendSignal{}
}
writeJSON(w, map[string]interface{}{"signals": signals})
}
func (d *Dashboard) handleSurgeEvents(w http.ResponseWriter, r *http.Request) {
limit := 100
// Try DB first
if d.db != nil {
events, err := d.db.GetSurgeEvents(limit)
if err == nil {
writeJSON(w, map[string]interface{}{"events": events, "total": len(events)})
return
}
}
// Fallback to in-memory
if d.surgeDetector != nil {
events := d.surgeDetector.GetRecentEvents(limit)
writeJSON(w, map[string]interface{}{"events": events, "total": len(events)})
} else {
writeJSON(w, map[string]interface{}{"events": []interface{}{}, "total": 0})
}
}
func (d *Dashboard) handleSSE(w http.ResponseWriter, r *http.Request) {
flusher, ok := w.(http.Flusher)
if !ok {
http.Error(w, "Streaming not supported", 500)
return
}
w.Header().Set("Content-Type", "text/event-stream")
w.Header().Set("Cache-Control", "no-cache")
w.Header().Set("Connection", "keep-alive")
w.Header().Set("Access-Control-Allow-Origin", "*")
client := d.hub.Subscribe("")
defer d.hub.Unsubscribe(client)
fmt.Fprintf(w, "event: connected\ndata: {\"status\":\"ok\"}\n\n")
flusher.Flush()
for {
select {
case <-r.Context().Done():
return
case msg, ok := <-client.ch:
if !ok {
return
}
fmt.Fprintf(w, "data: %s\n\n", msg)
flusher.Flush()
}
}
}
// handleBinanceIndex serves the Binance momentum monitor page.
func (d *Dashboard) handleBinanceIndex(w http.ResponseWriter, r *http.Request) {
var data []byte
var err error
data, err = os.ReadFile("frontend/dist/binance.html")
if err != nil {
data, err = staticFS.ReadFile("frontend/dist/binance.html")
}
if err != nil {
http.Error(w, "Not found", 404)
return
}
w.Header().Set("Content-Type", "text/html; charset=utf-8")
w.Write(data)
}
// handleBinanceAlerts returns recent Binance momentum alerts.
func (d *Dashboard) handleBinanceAlerts(w http.ResponseWriter, r *http.Request) {
if d.binanceMomentumDetector == nil {
writeJSON(w, map[string]interface{}{"alerts": []interface{}{}})
return
}
alerts := d.binanceMomentumDetector.RecentAlerts(50)
if alerts == nil {
alerts = []BinanceAlert{}
}
writeJSON(w, map[string]interface{}{"alerts": alerts})
}
func writeJSON(w http.ResponseWriter, v interface{}) {
w.Header().Set("Content-Type", "application/json")
json.NewEncoder(w).Encode(v)
}
func trackedCoinNames() []string {
names := make([]string, len(TrackedCoins))
for i, c := range TrackedCoins {
names[i] = c.Name
}
return names
}