Files
ryan b396638438 feat(交易信号): 持仓状态持久化至 Redis 以保障重启与多实例下均价计算准确
Motivation:
持仓与已处理信号的状态此前仅存于进程内存,服务重启或多实例部署后会丢失,导致加仓均价、仓位大小等计算失真,重复信号也无法跨实例幂等。通过将状态持久化到 Redis,保证持仓跟踪跨重启、跨实例连续一致,提升通知内容的准确性与可靠性。

Changes:

* 新增持仓存储抽象,支持内存与 Redis 两种实现,持仓状态与已处理信号快照按 TTL 持久化
* 持仓变更通过 Redis 事务管道原子提交,保证状态更新与幂等记录一致写入
* 存储写入失败时消息进入重试而非直接确认,避免状态丢失导致通知失真
* 缓存层新增原始值读取与批量事务写入能力,并在订阅器初始化时注入 Redis 依赖
* 补充跨实例持久化、幂等去重与存储失败场景的测试覆盖
2026-08-23 14:43:53 +08:00

101 lines
2.5 KiB
Go

package tradesignal
import (
"encoding/json"
"errors"
"fmt"
"strings"
"aiaa-notification-service/internal/cache"
"aiaa-notification-service/internal/config"
)
var ErrInvalidSignal = errors.New("invalid signal")
var ErrPositionStore = errors.New("position store")
type Converter struct {
overrides map[string]config.StrategyOverride
positions *Tracker
}
func NewConverter(overrides map[string]config.StrategyOverride) *Converter {
return NewConverterWithCache(overrides, nil)
}
func NewConverterWithCache(overrides map[string]config.StrategyOverride, c *cache.Cache) *Converter {
return &Converter{
overrides: overrides,
positions: NewTrackerWithStore(newRedisStore(c)),
}
}
func (c *Converter) Convert(body []byte) (string, map[string]interface{}, error) {
var sig Signal
if err := json.Unmarshal(body, &sig); err != nil {
return "", nil, fmt.Errorf("%w: %v", ErrInvalidSignal, err)
}
if strings.TrimSpace(sig.Action) == "" {
if strings.TrimSpace(sig.RawMessage) == "" {
return "", nil, fmt.Errorf("%w: missing action", ErrInvalidSignal)
}
data, err := toData(body, &sig)
if err != nil {
return "", nil, err
}
return "trade.message", data, nil
}
out := Apply(&sig, c.overrideFor(sig.StrategyCode))
snap, err := c.positions.Apply(out)
if err != nil {
return "", nil, fmt.Errorf("%w: %v", ErrPositionStore, err)
}
var opts FormatOptions
if snap.HasAvg {
avg := snap.AvgPrice
opts.AvgPrice = &avg
}
text := Format(out, opts)
data, err := toData(body, out)
if err != nil {
return "", nil, err
}
data["formatted"] = text
if snap.HasAvg {
data["avgPrice"] = snap.AvgPrice
}
return "trade." + strings.ToLower(out.Action), data, nil
}
func (c *Converter) overrideFor(code string) *config.StrategyOverride {
if c == nil || len(c.overrides) == 0 || code == "" {
return nil
}
// viper lower-cases nested map keys, so match case-insensitively.
if override, ok := c.overrides[code]; ok {
return &override
}
if override, ok := c.overrides[strings.ToLower(code)]; ok {
return &override
}
return nil
}
func toData(body []byte, sig *Signal) (map[string]interface{}, error) {
data := make(map[string]interface{})
if err := json.Unmarshal(body, &data); err != nil {
return nil, err
}
raw, err := json.Marshal(sig)
if err != nil {
return nil, err
}
overlay := make(map[string]interface{})
if err := json.Unmarshal(raw, &overlay); err != nil {
return nil, err
}
for k, v := range overlay {
data[k] = v
}
return data, nil
}