284 lines
8.2 KiB
Go
284 lines
8.2 KiB
Go
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)
|
||
}
|
||
}
|
||
}
|