189 lines
4.1 KiB
Go
189 lines
4.1 KiB
Go
package main
|
|
|
|
import (
|
|
"database/sql"
|
|
"encoding/json"
|
|
"fmt"
|
|
"log"
|
|
"strconv"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// LiquidationHandler handles allLiquidation stream messages, saving events to hot DB and bucketing 5s features.
|
|
type LiquidationHandler struct {
|
|
cfg Config
|
|
storage *StreamStorage
|
|
mu sync.Mutex
|
|
currentBucket int64
|
|
events []Liquidation
|
|
hotDB *sql.DB
|
|
featDB *sql.DB
|
|
}
|
|
|
|
func NewLiquidationHandler(cfg Config, sm *StorageManager) (*LiquidationHandler, error) {
|
|
ss := sm.GetStreamStorage("liquidations")
|
|
if ss == nil {
|
|
return nil, fmt.Errorf("liquidations stream storage not found")
|
|
}
|
|
|
|
hotDB, err := OpenDBWithAutoVacuum(ss.DBPath("hot_liquidations.db"))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("open hot_liquidations db: %w", err)
|
|
}
|
|
|
|
featDB, err := OpenDBWithAutoVacuum(ss.DBPath("features.db"))
|
|
if err != nil {
|
|
hotDB.Close()
|
|
return nil, fmt.Errorf("open liquidation features db: %w", err)
|
|
}
|
|
|
|
return &LiquidationHandler{
|
|
cfg: cfg,
|
|
storage: ss,
|
|
hotDB: hotDB,
|
|
featDB: featDB,
|
|
}, nil
|
|
}
|
|
|
|
func (lh *LiquidationHandler) Topics() []string {
|
|
return []string{fmt.Sprintf("allLiquidation.%s", lh.cfg.Symbol)}
|
|
}
|
|
|
|
func (lh *LiquidationHandler) HandleMessage(data []byte) {
|
|
var msg BybitLiquidationMessage
|
|
if err := json.Unmarshal(data, &msg); err != nil {
|
|
return
|
|
}
|
|
|
|
raw := msg.Data
|
|
if raw.S != "" && raw.S != lh.cfg.Symbol {
|
|
return // Filter by target symbol
|
|
}
|
|
|
|
price, err := strconv.ParseFloat(raw.P, 64)
|
|
if err != nil {
|
|
return
|
|
}
|
|
qty, err := strconv.ParseFloat(raw.V, 64)
|
|
if err != nil {
|
|
return
|
|
}
|
|
|
|
ts := raw.T
|
|
if ts == 0 {
|
|
ts = msg.TS
|
|
}
|
|
if ts == 0 {
|
|
ts = time.Now().UnixMilli()
|
|
}
|
|
|
|
liq := Liquidation{
|
|
Timestamp: ts,
|
|
Side: raw.SD,
|
|
Price: price,
|
|
Quantity: qty,
|
|
Value: price * qty,
|
|
}
|
|
|
|
// 1. Write event directly to hot DB
|
|
_, err = lh.hotDB.Exec(`
|
|
INSERT INTO liquidations (timestamp, side, price, quantity, value)
|
|
VALUES (?, ?, ?, ?, ?)
|
|
`, liq.Timestamp, liq.Side, liq.Price, liq.Quantity, liq.Value)
|
|
if err != nil {
|
|
log.Printf("[liquidation_handler] hot db insert error: %v", err)
|
|
}
|
|
|
|
// 2. Aggregate into 5-second feature buckets
|
|
lh.mu.Lock()
|
|
defer lh.mu.Unlock()
|
|
|
|
bucketTS := (liq.Timestamp / 5000) * 5000
|
|
|
|
if lh.currentBucket == 0 {
|
|
lh.currentBucket = bucketTS
|
|
}
|
|
|
|
if bucketTS > lh.currentBucket {
|
|
if len(lh.events) > 0 {
|
|
lh.flushBucket()
|
|
}
|
|
lh.currentBucket = bucketTS
|
|
lh.events = lh.events[:0]
|
|
}
|
|
|
|
lh.events = append(lh.events, liq)
|
|
}
|
|
|
|
func (lh *LiquidationHandler) flushBucket() {
|
|
if len(lh.events) == 0 {
|
|
return
|
|
}
|
|
|
|
feat := computeLiquidationFeatures(lh.currentBucket, lh.events)
|
|
|
|
_, err := lh.featDB.Exec(`
|
|
INSERT OR IGNORE INTO liquidation_features (
|
|
timestamp, count_total, count_long, count_short,
|
|
volume_total, volume_long, volume_short,
|
|
value_total, value_long, value_short,
|
|
avg_price, net_value
|
|
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
|
`, feat.Timestamp, feat.CountTotal, feat.CountLong, feat.CountShort,
|
|
feat.VolumeTotal, feat.VolumeLong, feat.VolumeShort,
|
|
feat.ValueTotal, feat.ValueLong, feat.ValueShort,
|
|
feat.AvgPrice, feat.NetValue)
|
|
|
|
if err != nil {
|
|
log.Printf("[liquidation_handler] feature insert error: %v", err)
|
|
}
|
|
}
|
|
|
|
func computeLiquidationFeatures(bucketTS int64, events []Liquidation) LiquidationFeature {
|
|
feat := LiquidationFeature{
|
|
Timestamp: bucketTS,
|
|
CountTotal: len(events),
|
|
}
|
|
|
|
var sumPrice float64
|
|
for _, ev := range events {
|
|
sumPrice += ev.Price
|
|
feat.VolumeTotal += ev.Quantity
|
|
feat.ValueTotal += ev.Value
|
|
|
|
if ev.Side == "Buy" { // Buy = long liquidated
|
|
feat.CountLong++
|
|
feat.VolumeLong += ev.Quantity
|
|
feat.ValueLong += ev.Value
|
|
} else { // Sell = short liquidated
|
|
feat.CountShort++
|
|
feat.VolumeShort += ev.Quantity
|
|
feat.ValueShort += ev.Value
|
|
}
|
|
}
|
|
|
|
if feat.CountTotal > 0 {
|
|
feat.AvgPrice = sumPrice / float64(feat.CountTotal)
|
|
}
|
|
feat.NetValue = feat.ValueLong - feat.ValueShort
|
|
|
|
return feat
|
|
}
|
|
|
|
func (lh *LiquidationHandler) Close() {
|
|
lh.mu.Lock()
|
|
defer lh.mu.Unlock()
|
|
|
|
if len(lh.events) > 0 {
|
|
lh.flushBucket()
|
|
}
|
|
|
|
if lh.hotDB != nil {
|
|
lh.hotDB.Close()
|
|
}
|
|
if lh.featDB != nil {
|
|
lh.featDB.Close()
|
|
}
|
|
}
|