Files

261 lines
6.1 KiB
Go

package main
import (
"context"
"fmt"
"log"
"os"
"os/exec"
"os/signal"
"sync"
"syscall"
"time"
"bybit_btcusdt_ingest/utils/stats"
)
func main() {
log.SetFlags(log.LstdFlags | log.Lmicroseconds)
// Load configuration
cfg, err := LoadConfig("config.json")
if err != nil {
log.Fatalf("Config error: %v", err)
}
if len(os.Args) < 2 {
printUsage()
return
}
command := os.Args[1]
switch command {
case "run":
isDaemon := false
for _, arg := range os.Args {
if arg == "--daemon" {
isDaemon = true
break
}
}
if isDaemon {
startDaemon(cfg.LogFile)
return
}
runEngine(cfg)
case "recover":
log.Println("=== Starting Recovery Mode ===")
sm, err := NewStorageManager(cfg)
if err != nil {
log.Fatal(err)
}
if err := sm.StartupRecovery(); err != nil {
log.Fatal(err)
}
log.Println("Recovery completed.")
case "maintain":
log.Println("=== Starting Manual Maintenance ===")
sm, err := NewStorageManager(cfg)
if err != nil {
log.Fatal(err)
}
sm.RunHourlyMaintenance()
log.Println("Maintenance completed.")
case "stats":
dataDir := cfg.DataDir
noColor := false
for i := 2; i < len(os.Args); i++ {
arg := os.Args[i]
if arg == "--no-color" {
noColor = true
} else if (arg == "--data" || arg == "-data") && i+1 < len(os.Args) {
dataDir = os.Args[i+1]
i++
}
}
stats.PrintDashboard(dataDir, noColor)
case "help":
printUsage()
default:
log.Printf("Unknown command: %s", command)
printUsage()
}
}
func printUsage() {
fmt.Println("Usage: engine [command] [--daemon]")
fmt.Println("\nCommands:")
fmt.Println(" run Start the Multi-Stream WebSocket ingestor and processing engine")
fmt.Println(" Use --daemon to run in background and log to the file defined in config")
fmt.Println(" recover Run startup recovery to migrate stale trade ticks")
fmt.Println(" maintain Run a single maintenance cycle (cleanup/retention across all streams)")
fmt.Println(" stats Display database statistics across all data streams")
fmt.Println(" Use --data <path> to specify custom data directory")
fmt.Println(" Use --no-color to disable ANSI color formatting")
fmt.Println(" help Show this help message")
}
func startDaemon(logPath string) {
if logPath == "" {
logPath = "engine.log"
}
args := []string{}
for _, arg := range os.Args[1:] {
if arg != "--daemon" {
args = append(args, arg)
}
}
cmd := exec.Command(os.Args[0], args...)
logFile, err := os.OpenFile(logPath, os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0666)
if err != nil {
log.Fatalf("Failed to open log file %s: %v", logPath, err)
}
cmd.Stdout = logFile
cmd.Stderr = logFile
if err := cmd.Start(); err != nil {
log.Fatalf("Failed to start daemon: %v", err)
}
fmt.Printf("Engine started in background (PID: %d). Logs: %s\n", cmd.Process.Pid, logPath)
os.Exit(0)
}
func runEngine(cfg Config) {
log.Println("=== Bybit Multi-Stream Ingest Engine ===")
log.Printf("Symbol: %s | Data Dir: %s", cfg.Symbol, cfg.DataDir)
// Initialize storage (creates per-stream directories, databases, tables, and migrates old layout if needed)
sm, err := NewStorageManager(cfg)
if err != nil {
log.Fatalf("Storage init error: %v", err)
}
// Startup recovery for trades
if err := sm.StartupRecovery(); err != nil {
log.Fatalf("Startup recovery error: %v", err)
}
handlers := make(map[string]MessageHandler)
// Initialize enabled handlers
if cfg.Streams.Trades.Enabled {
th, err := NewTradeHandler(cfg, sm)
if err != nil {
log.Fatalf("TradeHandler init error: %v", err)
}
handlers["publicTrade"] = th
log.Println("[init] Trades stream handler enabled.")
}
if cfg.Streams.Ticker.Enabled {
tickerH, err := NewTickerHandler(cfg, sm)
if err != nil {
log.Fatalf("TickerHandler init error: %v", err)
}
handlers["tickers"] = tickerH
log.Println("[init] Tickers stream handler enabled.")
}
if cfg.Streams.Klines.Enabled {
klineH, err := NewKlineHandler(cfg, sm)
if err != nil {
log.Fatalf("KlineHandler init error: %v", err)
}
handlers["kline"] = klineH
log.Printf("[init] Klines stream handler enabled (intervals: %v).", cfg.Streams.Klines.Intervals)
}
if cfg.Streams.Orderbook.Enabled {
obH, err := NewOrderbookHandler(cfg, sm)
if err != nil {
log.Fatalf("OrderbookHandler init error: %v", err)
}
handlers["orderbook"] = obH
log.Printf("[init] Orderbook stream handler enabled (depth: %d).", cfg.Streams.Orderbook.Depth)
}
if cfg.Streams.Liquidations.Enabled {
liqH, err := NewLiquidationHandler(cfg, sm)
if err != nil {
log.Fatalf("LiquidationHandler init error: %v", err)
}
handlers["allLiquidation"] = liqH
log.Println("[init] Liquidations stream handler enabled.")
}
// Create multi-stream WebSocket Ingestor
ingestor := NewIngestor(cfg, handlers)
// Context for graceful shutdown
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
sigCh := make(chan os.Signal, 1)
signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM)
var wg sync.WaitGroup
// Goroutine 1: Trade Handler background writer (if trades enabled)
if th, ok := handlers["publicTrade"].(*TradeHandler); ok {
wg.Add(1)
go func() {
defer wg.Done()
th.writer.Run(ctx)
}()
}
// Goroutine 2: WebSocket Ingestor
wg.Add(1)
go func() {
defer wg.Done()
ingestor.Run(ctx)
}()
// Goroutine 3: Periodic Maintenance
wg.Add(1)
go func() {
defer wg.Done()
maintInterval := time.Duration(cfg.MaintenanceIntervalMinutes) * time.Minute
ticker := time.NewTicker(maintInterval)
defer ticker.Stop()
log.Printf("[maintenance] Scheduled every %d minutes.", cfg.MaintenanceIntervalMinutes)
for {
select {
case <-ctx.Done():
log.Println("[maintenance] Maintenance goroutine stopped.")
return
case <-ticker.C:
sm.RunHourlyMaintenance()
}
}
}()
// Wait for shutdown signal
sig := <-sigCh
log.Printf("Received signal %v, initiating graceful shutdown...", sig)
cancel()
// Close all handlers
for name, h := range handlers {
log.Printf("[shutdown] Closing handler %s...", name)
h.Close()
}
wg.Wait()
log.Println("=== Engine shut down cleanly. ===")
}