package main import ( "database/sql" "encoding/json" "fmt" "log" "math" "strconv" "sync" ) // KlineHandler handles Bybit kline streams for multiple intervals (e.g., 5, 15, 60). type KlineHandler struct { cfg Config storage *StreamStorage dbs map[string]*sql.DB // interval -> kline_.db handle featDB *sql.DB mu sync.Mutex } func NewKlineHandler(cfg Config, sm *StorageManager) (*KlineHandler, error) { ss := sm.GetStreamStorage("klines") if ss == nil { return nil, fmt.Errorf("klines stream storage not found") } dbs := make(map[string]*sql.DB) for _, interval := range cfg.Streams.Klines.Intervals { dbName := fmt.Sprintf("kline_%s.db", interval) db, err := OpenDBWithAutoVacuum(ss.DBPath(dbName)) if err != nil { for _, d := range dbs { d.Close() } return nil, fmt.Errorf("open kline_%s db: %w", interval, err) } dbs[interval] = db } featDB, err := OpenDBWithAutoVacuum(ss.DBPath("features.db")) if err != nil { for _, d := range dbs { d.Close() } return nil, fmt.Errorf("open kline features db: %w", err) } return &KlineHandler{ cfg: cfg, storage: ss, dbs: dbs, featDB: featDB, }, nil } func (kh *KlineHandler) Topics() []string { var topics []string for _, interval := range kh.cfg.Streams.Klines.Intervals { topics = append(topics, fmt.Sprintf("kline.%s.%s", interval, kh.cfg.Symbol)) } return topics } func (kh *KlineHandler) HandleMessage(data []byte) { var msg BybitKlineMessage if err := json.Unmarshal(data, &msg); err != nil { return } if len(msg.Data) == 0 { return } kh.mu.Lock() defer kh.mu.Unlock() for _, raw := range msg.Data { open, _ := strconv.ParseFloat(raw.Open, 64) closeP, _ := strconv.ParseFloat(raw.Close, 64) high, _ := strconv.ParseFloat(raw.High, 64) low, _ := strconv.ParseFloat(raw.Low, 64) vol, _ := strconv.ParseFloat(raw.Volume, 64) turnover, _ := strconv.ParseFloat(raw.Turnover, 64) confirmInt := 0 if raw.Confirm { confirmInt = 1 } db, ok := kh.dbs[raw.Interval] if !ok { continue } // Insert or replace kline in hot DB _, err := db.Exec(` INSERT INTO klines (start_time, end_time, interval, open, high, low, close, volume, turnover, confirmed) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(start_time) DO UPDATE SET end_time=excluded.end_time, open=excluded.open, high=excluded.high, low=excluded.low, close=excluded.close, volume=excluded.volume, turnover=excluded.turnover, confirmed=excluded.confirmed `, raw.Start, raw.End, raw.Interval, open, high, low, closeP, vol, turnover, confirmInt) if err != nil { log.Printf("[kline_handler] insert kline error: %v", err) } // Calculate features if confirmed or on updates if raw.Confirm { highLow := high - low bodyRatio := 0.0 upperWick := 0.0 lowerWick := 0.0 if highLow > 0 { bodyRatio = math.Abs(closeP-open) / highLow maxBody := math.Max(open, closeP) minBody := math.Min(open, closeP) upperWick = (high - maxBody) / highLow lowerWick = (minBody - low) / highLow } logRet := 0.0 if open > 0 { logRet = math.Log(closeP / open) } _, err = kh.featDB.Exec(` INSERT INTO kline_features (timestamp, interval, body_ratio, upper_wick, lower_wick, log_return, volume, turnover) VALUES (?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(interval, timestamp) DO UPDATE SET body_ratio=excluded.body_ratio, upper_wick=excluded.upper_wick, lower_wick=excluded.lower_wick, log_return=excluded.log_return, volume=excluded.volume, turnover=excluded.turnover `, raw.Start, raw.Interval, bodyRatio, upperWick, lowerWick, logRet, vol, turnover) if err != nil { log.Printf("[kline_handler] insert kline features error: %v", err) } } } } func (kh *KlineHandler) Close() { kh.mu.Lock() defer kh.mu.Unlock() for _, db := range kh.dbs { db.Close() } if kh.featDB != nil { kh.featDB.Close() } }