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 } func NewDashboard(store *PriceStore, database *db.DB, addr string, cfg *Config, momentumTracker *MomentumTracker, trendDetector *TrendDetector, cumulativeTracker *CumulativeTracker, trendFilter *TrendFilter, surgeDetector *SurgeDetector) *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, } // 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) 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) } } } } // ============================================================ // 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() } } } 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 }