package main import ( "bytes" "context" "io" "log" "math/rand" "os" "os/signal" "strings" "syscall" "time" "exchange-monitor/db" "exchange-monitor/exchange" ) func main() { log.SetFlags(log.Ldate | log.Ltime | log.Lshortfile) // Set up multi-writer: stdout + log file logPath := os.ExpandEnv("$HOME/Project/exchange-monitor-go/exchange-monitor.log") logFile, err := os.OpenFile(logPath, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0644) if err == nil { multi := io.MultiWriter(os.Stdout, logFile) log.SetOutput(multi) } else { log.SetOutput(os.Stdout) } log.Println("[Exchange Monitor] Starting surge detection mode...") loadDotEnv() cfg := LoadConfig() // Initialize Telegram sender telegramSender := NewTelegramSender(cfg.TelegramBotToken, cfg.TelegramChatID) if telegramSender.IsEnabled() { log.Printf("[Telegram] Alerts enabled -> chat %s", cfg.TelegramChatID) } else { log.Println("[Telegram] Alerts disabled (missing token or chat ID)") } store := NewPriceStore() // Initialize momentum tracker (for momentum scanning mode) momentumTracker := NewMomentumTracker() // Initialize trend detector (for price anomaly / trend detection) trendDetector := NewTrendDetector(momentumTracker) if cfg.TrendEnabled { trendDetector.Configure(cfg.TrendBaselineWindow, cfg.TrendAnomalyMul, cfg.TrendConfirmTicks, cfg.TrendAlertCooldown) log.Printf("[Trend] Z-score detection enabled (z-score >= %.1fσ, window=%d ticks, confirm=%d ticks)", cfg.TrendAnomalyMul, cfg.TrendBaselineWindow, cfg.TrendConfirmTicks) } // Initialize cumulative tracker (1min/5min multi-exchange consensus change) cumulativeTracker := NewCumulativeTracker() // Initialize trend filter (Binance K-line based quiet + EMA52 filter) trendFilter := NewTrendFilter(store, trendDetector) trendFilter.Start() defer trendFilter.Stop() // Initialize SQLite database database, err := db.Open("") if err != nil { log.Printf("[DB] Failed to open database: %v", err) } else { defer database.Close() } // Initialize surge detector surgeDetector := NewSurgeDetector() if cfg.SurgeEnabled { surgeDetector.Configure(cfg.SurgeWindowSize, cfg.SurgeBaselineMultiplier, cfg.SurgeMinAbsSpreadPct, cfg.SurgeCooldownSec) log.Printf("[Surge] Adaptive detection enabled (window=%d ticks, multiplier=%.1fx, min_spread=%.2f%%, cooldown=%ds)", cfg.SurgeWindowSize, cfg.SurgeBaselineMultiplier, cfg.SurgeMinAbsSpreadPct, cfg.SurgeCooldownSec) // Wire surge event persistence to SQLite if database != nil { surgeDetector.SetOnEvent(func(ev SurgeEvent) { database.InsertSurgeEvent(ev.Coin, ev.Timestamp, ev.BnPrice, ev.OkxPrice, ev.BgPrice, ev.SpreadPct, ev.BaselinePct, ev.ThresholdPct, ev.Ratio, ev.Direction, ev.LeadingExchange, ev.MidPrice) }) } } // Initialize Binance momentum detector (1-minute Binance price change) binanceMomentumDetector := NewBinanceMomentumDetector( cfg.BinanceMomentumWindowSec, 50, // actual loop tick is fixed 50ms cfg.BinanceMomentumThresholdPct, cfg.BinanceMomentumCooldownSec, ) if cfg.BinanceMomentumEnabled { log.Printf("[BinanceMomentum] Detection enabled (window=%ds, threshold=%.1f%%, cooldown=%ds)", cfg.BinanceMomentumWindowSec, cfg.BinanceMomentumThresholdPct, cfg.BinanceMomentumCooldownSec) } // Initialize dashboard (web server + SSE) dashboard := NewDashboard(store, database, ":8888", cfg, momentumTracker, trendDetector, cumulativeTracker, trendFilter, surgeDetector, binanceMomentumDetector) go dashboard.Run() // Context for graceful shutdown ctx, cancel := context.WithCancel(context.Background()) defer cancel() sigCh := make(chan os.Signal, 1) signal.Notify(sigCh, os.Interrupt, syscall.SIGUSR1) // Collect symbols for all exchanges var bgSymbols, bnSymbols, okxSymbols []string for _, c := range TrackedCoins { if c.BG != "" { bgSymbols = append(bgSymbols, c.BG) } if c.BN != "" { bnSymbols = append(bnSymbols, c.BN) } if c.OK != "" { okxSymbols = append(okxSymbols, c.OK) } } // Start exchange WS connections startExchange := func(name string, runner func(func(string, float64, float64, float64)) error) { go func() { for { err := runner(func(coin string, price, bid, ask float64) { store.SetWithSpread(coin, name, price, bid, ask) dashboard.RecordPrice(coin, name, price) dashboard.RecordConnStatus(name) }) log.Printf("[%s] WS error: %v (reconnecting...)", name, err) select { case <-ctx.Done(): return case <-time.After(3 * time.Second): } } }() } startExchange("Bitget", exchange.NewBitgetWS(bgSymbols).Run) startExchange("Binance", exchange.NewBinanceWS(bnSymbols).Run) startExchange("OKX", exchange.NewOKXWS(okxSymbols).Run) log.Println("[Monitor] Waiting for initial data...") time.Sleep(10 * time.Second) // Main loop — fixed 50ms scan interval jitterMin, jitterMax := 50, 50 randInterval := func() time.Duration { return time.Duration(jitterMin+rand.Intn(jitterMax-jitterMin+1)) * time.Millisecond } scannerTick := time.NewTimer(randInterval()) statusTick := time.NewTicker(30 * time.Second) log.Printf("[Monitor] Scanner running every %dms", jitterMin) runLoop := true for runLoop { select { case sig := <-sigCh: if sig == syscall.SIGUSR1 { log.Printf("[Monitor] SIGUSR1 received — stats dump") continue } log.Println("[Monitor] Shutting down...") cancel() runLoop = false case <-statusTick.C: snap := store.GetAll() count := 0 for _, exMap := range snap { count += len(exMap) } log.Printf("[Status] %d prices / %d coins connected", count, len(snap)) // Show surge events in last 30s events := surgeDetector.GetRecentEvents(3) for _, ev := range events { if time.Since(ev.Timestamp) < 30*time.Second { log.Printf(" [Surge] %s %s spread=%.4f%% leading=%s", ev.Coin, ev.Direction, ev.SpreadPct, ev.LeadingExchange) } } case <-scannerTick.C: now := time.Now() snap := store.GetAll() // Feed Binance prices to momentum detector if cfg.BinanceMomentumEnabled { for _, tc := range TrackedCoins { if exMap := snap[tc.Name]; exMap != nil { if bnPrice, ok := exMap[ExBinance]; ok && bnPrice > 0 { binanceMomentumDetector.Record(tc.Name, bnPrice) } } } alerts := binanceMomentumDetector.Detect() for _, alert := range alerts { log.Printf("[BinanceMomentum] %s %s %.2f%% ($%.4f)", alert.Coin, alert.Direction, alert.ChangePct, alert.Price) dashboard.BroadcastEvent("binance_alert", alert) if telegramSender.IsEnabled() { go func(a BinanceAlert) { if err := telegramSender.SendAlert(a); err != nil { log.Printf("[Telegram] Failed to send alert for %s: %v", a.Coin, err) } }(alert) } } } // Feed prices to momentum tracker (for momentum scanning or trend detection) if cfg.MomentumEnabled || cfg.TrendEnabled { for coin, exMap := range snap { for ex, price := range exMap { momentumTracker.Record(coin, ex, price) } } } // Feed snapshots to cumulative tracker (always on) for _, tc := range TrackedCoins { exMap := snap[tc.Name] if exMap == nil || len(exMap) < 3 { continue } cumulativeTracker.Record(tc.Name, exMap) } // Run 3-exchange spread scan spreads := Scan3Ex(snap) dashboard.UpdateScan(spreads) // Run surge detection if cfg.SurgeEnabled { newEvents := surgeDetector.Tick(snap) for _, ev := range newEvents { dashboard.BroadcastEvent("surge_event", ev) } } scannerTick.Reset(randInterval()) _ = now } } log.Println("[Monitor] Stopped.") } func loadDotEnv() { envPath := os.ExpandEnv("$HOME/Project/exchange-monitor-go/.env") if _, err := os.Stat(envPath); err != nil { return } data, err := os.ReadFile(envPath) if err != nil { return } for _, line := range bytes.Split(data, []byte("\n")) { line = bytes.TrimSpace(line) if len(line) == 0 || line[0] == '#' { continue } parts := bytes.SplitN(line, []byte("="), 2) if len(parts) != 2 { continue } key := string(bytes.TrimSpace(parts[0])) val := string(bytes.TrimSpace(parts[1])) // Strip inline comments if idx := strings.Index(val, "#"); idx >= 0 { val = strings.TrimSpace(val[:idx]) } if os.Getenv(key) == "" { os.Setenv(key, val) } } }