mirror of
https://github.com/c9s/bbgo.git
synced 2024-11-22 06:53:52 +00:00
indicator: implement Subscribe method on PriceStream
This commit is contained in:
parent
775ad7d906
commit
ea1025d790
|
@ -8,7 +8,7 @@ import (
|
|||
"github.com/c9s/bbgo/pkg/types"
|
||||
)
|
||||
|
||||
func TestIndicatorSet_closeCache(t *testing.T) {
|
||||
func newTestIndicatorSet() *IndicatorSet {
|
||||
symbol := "BTCUSDT"
|
||||
store := NewMarketDataStore(symbol)
|
||||
store.KLineWindows[types.Interval1m] = &types.KLineWindow{
|
||||
|
@ -17,15 +17,32 @@ func TestIndicatorSet_closeCache(t *testing.T) {
|
|||
{Open: number(19200.0), Close: number(19300.0)},
|
||||
{Open: number(19300.0), Close: number(19200.0)},
|
||||
{Open: number(19200.0), Close: number(19100.0)},
|
||||
{Open: number(19100.0), Close: number(19500.0)},
|
||||
{Open: number(19500.0), Close: number(19600.0)},
|
||||
{Open: number(19600.0), Close: number(19700.0)},
|
||||
}
|
||||
|
||||
stream := types.NewStandardStream()
|
||||
|
||||
indicatorSet := NewIndicatorSet(symbol, &stream, store)
|
||||
|
||||
return indicatorSet
|
||||
}
|
||||
|
||||
func TestIndicatorSet_closeCache(t *testing.T) {
|
||||
indicatorSet := newTestIndicatorSet()
|
||||
|
||||
close1m := indicatorSet.CLOSE(types.Interval1m)
|
||||
assert.NotNil(t, close1m)
|
||||
|
||||
close1m2 := indicatorSet.CLOSE(types.Interval1m)
|
||||
assert.Equal(t, close1m, close1m2)
|
||||
}
|
||||
|
||||
func TestIndicatorSet_rsi(t *testing.T) {
|
||||
indicatorSet := newTestIndicatorSet()
|
||||
|
||||
rsi1m := indicatorSet.RSI(types.IntervalWindow{Interval: types.Interval1m, Window: 7})
|
||||
assert.NotNil(t, rsi1m)
|
||||
|
||||
rsiLast := rsi1m.Last(0)
|
||||
assert.InDelta(t, 80, rsiLast, 0.0000001)
|
||||
}
|
||||
|
|
|
@ -22,15 +22,31 @@ func Price(source KLineSubscription, mapper KLineValueMapper) *PriceStream {
|
|||
mapper: mapper,
|
||||
}
|
||||
|
||||
if source != nil {
|
||||
source.AddSubscriber(func(k types.KLine) {
|
||||
v := s.mapper(k)
|
||||
s.PushAndEmit(v)
|
||||
})
|
||||
if source == nil {
|
||||
return s
|
||||
}
|
||||
|
||||
source.AddSubscriber(func(k types.KLine) {
|
||||
v := s.mapper(k)
|
||||
s.PushAndEmit(v)
|
||||
})
|
||||
return s
|
||||
}
|
||||
|
||||
// AddSubscriber adds the subscriber function and push historical data to the subscriber
|
||||
func (s *PriceStream) AddSubscriber(f func(v float64)) {
|
||||
s.OnUpdate(f)
|
||||
|
||||
if len(s.slice) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
// push historical value to the subscriber
|
||||
for _, v := range s.slice {
|
||||
f(v)
|
||||
}
|
||||
}
|
||||
|
||||
func (s *PriceStream) PushAndEmit(v float64) {
|
||||
s.slice.Push(v)
|
||||
s.EmitUpdate(v)
|
||||
|
|
Loading…
Reference in New Issue
Block a user