- 新增OKX WebSocket行情连接器,扩展4交易所价格监控 - 新增z-score趋势检测引擎(TrendDetector),识别价格异动/趋势启动 - 新增累积变动跟踪(CumulativeTracker),基于1min/5min多交易所共识 - 趋势事件和累积变动事件持久化到SQLite - 新增Binance/OKX动量检测字段,扩展前端动量卡片至15列 - 迁移至macOS(darwin-arm64),更新前端依赖 - Dashboard网格重构:非交易卡片置顶,交易卡片置底 - TrackedCoin添加OK字段,添加ExBinance/ExOKX常量 - 前端新增趋势检测卡片、趋势历史卡片、累积变动卡片 Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
325 lines
10 KiB
Go
325 lines
10 KiB
Go
package main
|
||
|
||
import (
|
||
"bytes"
|
||
"context"
|
||
"fmt"
|
||
"io"
|
||
"log"
|
||
"math/rand"
|
||
"os"
|
||
"os/signal"
|
||
"strings"
|
||
"syscall"
|
||
"time"
|
||
|
||
"exchange-monitor/db"
|
||
"exchange-monitor/exchange"
|
||
)
|
||
|
||
func main() {
|
||
// CLI subcommand mode: talk to running daemon via IPC
|
||
if len(os.Args) > 1 {
|
||
switch os.Args[1] {
|
||
case "status", "close-all", "stop", "start":
|
||
runIPCClient(os.Args[1], "")
|
||
case "close":
|
||
if len(os.Args) < 3 {
|
||
fmt.Fprintln(os.Stderr, "Usage: exchange-monitor close <COIN>")
|
||
os.Exit(1)
|
||
}
|
||
runIPCClient("close", os.Args[2])
|
||
default:
|
||
fmt.Fprintf(os.Stderr, "Unknown command: %s\n", os.Args[1])
|
||
fmt.Fprintln(os.Stderr, "Commands: status, close-all, close <COIN>, stop, start")
|
||
os.Exit(1)
|
||
}
|
||
return
|
||
}
|
||
|
||
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...")
|
||
|
||
loadDotEnv()
|
||
cfg := LoadConfig()
|
||
|
||
// Populate package-level taker fees from config (so scanner/dashboard/trader all use it)
|
||
takerFees[ExBitget] = cfg.TakerFeeBitget
|
||
takerFees[ExHyperLiquid] = cfg.TakerFeeHyperLiquid
|
||
|
||
store := NewPriceStore()
|
||
notifier := NewNotifier(cfg.TelegramBotToken, cfg.TelegramChatID)
|
||
|
||
// 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()
|
||
log.Printf("[CM] Cumulative change tracking enabled (1m >= %.1f%%, 3+ exchanges)", cumulativeTracker.surgePct1m)
|
||
|
||
// Initialize SQLite database
|
||
database, err := db.Open("")
|
||
if err != nil {
|
||
log.Printf("[DB] Failed to open database: %v", err)
|
||
} else {
|
||
defer database.Close()
|
||
}
|
||
|
||
// Initialize trader
|
||
trader := NewTrader(cfg, database)
|
||
|
||
// Start Unix socket IPC for CLI commands
|
||
trader.startIPCServer()
|
||
|
||
// Initialize dashboard (web server + SSE)
|
||
dashboard := NewDashboard(store, trader, database, ":8888", cfg, momentumTracker, trendDetector, cumulativeTracker)
|
||
go dashboard.Run()
|
||
|
||
// Spread window tracker — measures how long spreads stay above threshold
|
||
spreadTracker := NewSpreadWindowTracker()
|
||
|
||
// P3-4: wire real-time trade event broadcast
|
||
trader.OnTradeEvent = dashboard.BroadcastEvent
|
||
if cfg.MomentumEnabled {
|
||
log.Printf("[Trader] MOMENTUM SCAN mode: arbitrage trading disabled, momentum detection active (threshold >= %.2f%%)", cfg.MomentumThresholdPct)
|
||
} else if trader.IsConfigured() {
|
||
log.Printf("[Trader] %s mode: automated trading ENABLED (threshold >= %.2f%%, $%.0f/leg, max %d positions, $%.0f capital)",
|
||
trader.ModeLabel(), cfg.TradeThreshold, cfg.TradeAmountUSD, cfg.MaxPositions, cfg.InitialCapital)
|
||
if cfg.TestMode {
|
||
log.Printf("[Trader] Using mock orders with %.3f%% slippage per leg", cfg.MockSlippagePct)
|
||
}
|
||
log.Printf("[Trader] Bitget+HL: BG->HL / HL->BG only")
|
||
} else {
|
||
log.Printf("[Trader] Automated trading DISABLED (set TRADE_ENABLED=1 or TEST_MODE=true in .env)")
|
||
}
|
||
|
||
// Context for graceful shutdown — replaces shared sigCh (B#1)
|
||
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, hlSymbols, bnSymbols, okxSymbols []string
|
||
for _, c := range TrackedCoins {
|
||
if c.BG != "" {
|
||
bgSymbols = append(bgSymbols, c.BG)
|
||
}
|
||
if c.HL != "" {
|
||
hlSymbols = append(hlSymbols, c.HL)
|
||
}
|
||
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) // P3-5
|
||
})
|
||
log.Printf("[%s] WS error: %v (reconnecting...)", name, err)
|
||
select {
|
||
case <-ctx.Done():
|
||
return
|
||
case <-time.After(3 * time.Second):
|
||
}
|
||
}
|
||
}()
|
||
}
|
||
|
||
startExchange("HyperLiquid", exchange.NewHyperLiquidWS(hlSymbols).Run)
|
||
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
|
||
lastHour := -1
|
||
|
||
// 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 {
|
||
// Dump stats on request
|
||
converged, diverged, flat, total := trader.GetClosedStats()
|
||
stats := fmt.Sprintf("=== 收敛统计 === %s\n", time.Now().Format("2006-01-02 15:04"))
|
||
stats += fmt.Sprintf(" 总交易数: %d\n", total)
|
||
stats += fmt.Sprintf(" 价差收敛: %d\n", converged)
|
||
stats += fmt.Sprintf(" 价差持平: %d\n", flat)
|
||
stats += fmt.Sprintf(" 价差发散: %d\n", diverged)
|
||
if total > 0 {
|
||
stats += fmt.Sprintf(" 收敛率: %.1f%%\n", float64(converged)/float64(total)*100)
|
||
}
|
||
log.Printf("[Monitor] SIGUSR1 received — wrote stats to trade_stats.txt")
|
||
statsPath := os.ExpandEnv("$HOME/Project/exchange-monitor-go/trade_stats.txt")
|
||
os.WriteFile(statsPath, []byte(stats), 0644)
|
||
continue
|
||
}
|
||
log.Println("[Monitor] Shutting down...")
|
||
cancel() // B#1: cancel context to stop all WS goroutines
|
||
runLoop = false
|
||
|
||
case <-trader.StopCh:
|
||
log.Println("[Monitor] 5 real trades completed — trading stopped. System still running (dashboard active)")
|
||
log.Println("[Monitor] Use POST /api/start to resume trading, POST /api/stop to stop manually")
|
||
|
||
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 open positions (read from decoupled snapshot)
|
||
if positions := trader.ReadSnapshot(); len(positions) > 0 {
|
||
for _, pos := range positions {
|
||
log.Printf(" [Position] %s %s open %d scales $%.0f since %s",
|
||
pos.Coin, pos.Direction, pos.ScaleLevels, pos.AmountUSD,
|
||
time.Since(pos.StartedAt).Round(time.Second).String())
|
||
}
|
||
}
|
||
|
||
case <-scannerTick.C:
|
||
now := time.Now()
|
||
t0 := now
|
||
|
||
// Tick the trader (monitor open positions for exit)
|
||
trader.Tick(store, notifier)
|
||
trader.RefreshSnapshot() // decoupled snapshot for display
|
||
t1 := time.Now()
|
||
|
||
// Scan for arbitrage entries using maker fees (limit orders)
|
||
snap := store.GetAll()
|
||
|
||
// 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)
|
||
}
|
||
|
||
makerOpps := ScanBGHL(snap)
|
||
dashboard.UpdateScan(makerOpps)
|
||
t2 := time.Now()
|
||
|
||
// Track spread window durations (how long each opportunity stays alive)
|
||
spreadTracker.Tick(snap, cfg.TradeThreshold)
|
||
|
||
// In momentum mode, arbitrage trading is disabled
|
||
if !cfg.MomentumEnabled {
|
||
for _, opp := range makerOpps {
|
||
if opp.NetProfit < cfg.ArbThreshold {
|
||
continue
|
||
}
|
||
if trader.TryEntry(opp, store, notifier) {
|
||
log.Printf("[Trader] %s: entry initiated for %.4f%%", opp.Coin, opp.NetProfit)
|
||
}
|
||
}
|
||
}
|
||
t3 := time.Now()
|
||
|
||
// Profile: warn if any step is slow
|
||
tickDur := t3.Sub(t0)
|
||
tickMs := tickDur.Milliseconds()
|
||
if tickMs > 100 || t1.Sub(t0) > 50*time.Millisecond || t2.Sub(t1) > 50*time.Millisecond || t3.Sub(t2) > 50*time.Millisecond {
|
||
log.Printf("[Profile] tick=%dms trader=%dms scan=%dms entry=%dms",
|
||
tickMs, t1.Sub(t0).Milliseconds(), t2.Sub(t1).Milliseconds(), t3.Sub(t2).Milliseconds())
|
||
}
|
||
|
||
// Hourly trade summary — use hour-based tracking (wider window than second-granularity)
|
||
hour := now.Hour()
|
||
if hour != lastHour && now.Minute() < 1 {
|
||
positions := trader.ReadSnapshot()
|
||
notifier.SendTradeSummary(positions, now.Format("2006-01-02 15:04"))
|
||
lastHour = hour
|
||
}
|
||
|
||
scannerTick.Reset(randInterval())
|
||
}
|
||
}
|
||
|
||
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)
|
||
}
|
||
}
|
||
}
|