Files
aiaa-notification-server/internal/subscriber/tradesignal/position.go
T
ryan 6f846a0a3c feat: subscribe to RabbitMQ trade signals and notify by rules
Consume configurable queues, format signals (including period), share NotifyService with HTTP, and drop duplicate bodies within 1h.
2026-08-15 17:34:49 +08:00

205 lines
3.7 KiB
Go

package tradesignal
import (
"strings"
"sync"
)
type mode int
const (
modeNone mode = iota
modeQty
modeWeight
)
type Snapshot struct {
AvgPrice float64
Size float64
HasAvg bool
}
type state struct {
avg float64
size float64
mode mode
}
type Tracker struct {
mu sync.Mutex
positions map[string]*state
applied map[string]Snapshot
}
func NewTracker() *Tracker {
return &Tracker{
positions: make(map[string]*state),
applied: make(map[string]Snapshot),
}
}
func (t *Tracker) Apply(signal *Signal) Snapshot {
if signal == nil {
return Snapshot{}
}
t.mu.Lock()
defer t.mu.Unlock()
if signal.SignalID != "" {
if snap, ok := t.applied[signal.SignalID]; ok {
return snap
}
}
key := positionKey(signal.StrategyCode, signal.Symbol, signal.Side)
action := strings.ToUpper(signal.Action)
st := t.positions[key]
var snap Snapshot
switch action {
case "OPEN":
st = openPosition(signal)
snap = snapshotFrom(st)
if st != nil {
t.positions[key] = st
} else {
delete(t.positions, key)
}
case "ADD":
st = addPosition(st, signal)
snap = snapshotFrom(st)
if st != nil {
t.positions[key] = st
}
case "REDUCE":
snap = snapshotFrom(st)
st = reducePosition(st, signal)
if st == nil || st.size <= 0 {
delete(t.positions, key)
} else {
t.positions[key] = st
}
case "CLOSE":
snap = snapshotFrom(st)
delete(t.positions, key)
default:
snap = snapshotFrom(st)
}
if signal.SignalID != "" {
t.applied[signal.SignalID] = snap
}
return snap
}
func openPosition(signal *Signal) *state {
if qty, ok := positiveQty(signal.Quantity); ok {
return &state{avg: signal.Price, size: qty, mode: modeQty}
}
if w, ok := positiveRatio(signal.AmountMarginRatio); ok {
return &state{avg: signal.Price, size: w, mode: modeWeight}
}
if signal.Price > 0 {
return &state{avg: signal.Price, size: 0, mode: modeNone}
}
return nil
}
func addPosition(st *state, signal *Signal) *state {
if st == nil || st.size <= 0 {
return openPosition(signal)
}
if qty, ok := positiveQty(signal.Quantity); ok {
if st.mode == modeWeight {
return st
}
if st.mode == modeNone || st.size == 0 {
st.mode = modeQty
st.size = qty
st.avg = signal.Price
return st
}
st.avg = (st.size*st.avg + qty*signal.Price) / (st.size + qty)
st.size += qty
st.mode = modeQty
return st
}
if w, ok := positiveRatio(signal.AmountMarginRatio); ok {
if st.mode == modeQty {
return st
}
if st.mode == modeNone || st.size == 0 {
st.mode = modeWeight
st.size = w
st.avg = signal.Price
return st
}
st.avg = (st.size*st.avg + w*signal.Price) / (st.size + w)
st.size += w
st.mode = modeWeight
return st
}
return st
}
func reducePosition(st *state, signal *Signal) *state {
if st == nil {
return nil
}
if qty, ok := positiveQty(signal.Quantity); ok && st.mode == modeQty {
st.size -= qty
if st.size < 0 {
st.size = 0
}
return st
}
ratio := 0.0
if r, ok := positiveRatio(signal.PosMarginRatio); ok {
ratio = r
} else if signal.Quantity == nil && signal.PosMarginRatio == nil {
return st
}
if ratio > 1 {
ratio = 1
}
if ratio > 0 {
st.size *= (1 - ratio)
}
return st
}
func snapshotFrom(st *state) Snapshot {
if st == nil || st.avg <= 0 {
return Snapshot{}
}
return Snapshot{
AvgPrice: st.avg,
Size: st.size,
HasAvg: true,
}
}
func positiveQty(q *float64) (float64, bool) {
if q == nil || *q <= 0 {
return 0, false
}
return *q, true
}
func positiveRatio(r *float64) (float64, bool) {
if r == nil || *r <= 0 {
return 0, false
}
return *r, true
}
func positionKey(strategyCode, symbol, side string) string {
return strings.ToUpper(strategyCode) + "|" + strings.ToUpper(symbol) + "|" + strings.ToUpper(side)
}