Testdata: introduce basic simulation framework (#47863)
This commit is contained in:
@@ -15,22 +15,18 @@ import (
|
||||
|
||||
var random20HzStreamRegex = regexp.MustCompile(`random-20Hz-stream(-\d+)?`)
|
||||
|
||||
func (s *Service) SubscribeStream(_ context.Context, req *backend.SubscribeStreamRequest) (*backend.SubscribeStreamResponse, error) {
|
||||
func (s *Service) SubscribeStream(ctx context.Context, req *backend.SubscribeStreamRequest) (*backend.SubscribeStreamResponse, error) {
|
||||
s.logger.Debug("Allowing access to stream", "path", req.Path, "user", req.PluginContext.User)
|
||||
|
||||
if strings.HasPrefix(req.Path, "sim/") {
|
||||
return s.sims.SubscribeStream(ctx, req)
|
||||
}
|
||||
|
||||
initialData, err := backend.NewInitialFrame(s.frame, data.IncludeSchemaOnly)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// For flight simulations, send the more complex schema
|
||||
if strings.HasPrefix(req.Path, "flight") {
|
||||
ff := newFlightConfig().initFields()
|
||||
initialData, err = backend.NewInitialFrame(ff.frame, data.IncludeSchemaOnly)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
if strings.Contains(req.Path, "-labeled") {
|
||||
initialData, err = backend.NewInitialFrame(s.labelFrame, data.IncludeSchemaOnly)
|
||||
if err != nil {
|
||||
@@ -49,8 +45,13 @@ func (s *Service) SubscribeStream(_ context.Context, req *backend.SubscribeStrea
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (s *Service) PublishStream(_ context.Context, req *backend.PublishStreamRequest) (*backend.PublishStreamResponse, error) {
|
||||
func (s *Service) PublishStream(ctx context.Context, req *backend.PublishStreamRequest) (*backend.PublishStreamResponse, error) {
|
||||
s.logger.Debug("Attempt to publish into stream", "path", req.Path, "user", req.PluginContext.User)
|
||||
|
||||
if strings.HasPrefix(req.Path, "sim/") {
|
||||
return s.sims.PublishStream(ctx, req)
|
||||
}
|
||||
|
||||
return &backend.PublishStreamResponse{
|
||||
Status: backend.PublishStreamStatusPermissionDenied,
|
||||
}, nil
|
||||
@@ -58,6 +59,11 @@ func (s *Service) PublishStream(_ context.Context, req *backend.PublishStreamReq
|
||||
|
||||
func (s *Service) RunStream(ctx context.Context, request *backend.RunStreamRequest, sender *backend.StreamSender) error {
|
||||
s.logger.Debug("New stream call", "path", request.Path)
|
||||
|
||||
if strings.HasPrefix(request.Path, "sim/") {
|
||||
return s.sims.RunStream(ctx, request, sender)
|
||||
}
|
||||
|
||||
var conf testStreamConfig
|
||||
switch {
|
||||
case request.Path == "random-2s-stream":
|
||||
@@ -75,11 +81,6 @@ func (s *Service) RunStream(ctx context.Context, request *backend.RunStreamReque
|
||||
Drop: 0.2, // keep 80%
|
||||
Labeled: true,
|
||||
}
|
||||
case request.Path == "flight-5hz-stream":
|
||||
conf = testStreamConfig{
|
||||
Interval: 200 * time.Millisecond,
|
||||
Flight: newFlightConfig(),
|
||||
}
|
||||
case random20HzStreamRegex.MatchString(request.Path):
|
||||
conf = testStreamConfig{
|
||||
Interval: 50 * time.Millisecond,
|
||||
@@ -93,7 +94,6 @@ func (s *Service) RunStream(ctx context.Context, request *backend.RunStreamReque
|
||||
type testStreamConfig struct {
|
||||
Interval time.Duration
|
||||
Drop float64
|
||||
Flight *flightConfig
|
||||
Labeled bool
|
||||
}
|
||||
|
||||
@@ -104,12 +104,6 @@ func (s *Service) runTestStream(ctx context.Context, path string, conf testStrea
|
||||
ticker := time.NewTicker(conf.Interval)
|
||||
defer ticker.Stop()
|
||||
|
||||
var flight *flightFields
|
||||
if conf.Flight != nil {
|
||||
flight = conf.Flight.initFields()
|
||||
flight.append(conf.Flight.getNextPoint(time.Now()))
|
||||
}
|
||||
|
||||
labelFrame := data.NewFrame("labeled",
|
||||
data.NewField("labels", nil, make([]string, 2)),
|
||||
data.NewField("Time", nil, make([]time.Time, 2)),
|
||||
@@ -131,37 +125,30 @@ func (s *Service) runTestStream(ctx context.Context, path string, conf testStrea
|
||||
mode = data.IncludeAll
|
||||
}
|
||||
|
||||
if flight != nil {
|
||||
flight.set(0, conf.Flight.getNextPoint(t))
|
||||
if err := sender.SendFrame(flight.frame, mode); err != nil {
|
||||
delta := rand.Float64() - 0.5
|
||||
walker += delta
|
||||
|
||||
if conf.Labeled {
|
||||
secA := t.Second() / 3
|
||||
secB := t.Second() / 7
|
||||
|
||||
labelFrame.Fields[0].Set(0, fmt.Sprintf("s=A,s=p%d,x=X", secA))
|
||||
labelFrame.Fields[1].Set(0, t)
|
||||
labelFrame.Fields[2].Set(0, walker)
|
||||
|
||||
labelFrame.Fields[0].Set(1, fmt.Sprintf("s=B,s=p%d,x=X", secB))
|
||||
labelFrame.Fields[1].Set(1, t)
|
||||
labelFrame.Fields[2].Set(1, walker+10)
|
||||
if err := sender.SendFrame(labelFrame, mode); err != nil {
|
||||
return err
|
||||
}
|
||||
} else {
|
||||
delta := rand.Float64() - 0.5
|
||||
walker += delta
|
||||
|
||||
if conf.Labeled {
|
||||
secA := t.Second() / 3
|
||||
secB := t.Second() / 7
|
||||
|
||||
labelFrame.Fields[0].Set(0, fmt.Sprintf("s=A,s=p%d,x=X", secA))
|
||||
labelFrame.Fields[1].Set(0, t)
|
||||
labelFrame.Fields[2].Set(0, walker)
|
||||
|
||||
labelFrame.Fields[0].Set(1, fmt.Sprintf("s=B,s=p%d,x=X", secB))
|
||||
labelFrame.Fields[1].Set(1, t)
|
||||
labelFrame.Fields[2].Set(1, walker+10)
|
||||
if err := sender.SendFrame(labelFrame, mode); err != nil {
|
||||
return err
|
||||
}
|
||||
} else {
|
||||
s.frame.Fields[0].Set(0, t)
|
||||
s.frame.Fields[1].Set(0, walker) // Value
|
||||
s.frame.Fields[2].Set(0, walker-((rand.Float64()*spread)+0.01)) // Min
|
||||
s.frame.Fields[3].Set(0, walker+((rand.Float64()*spread)+0.01)) // Max
|
||||
if err := sender.SendFrame(s.frame, mode); err != nil {
|
||||
return err
|
||||
}
|
||||
s.frame.Fields[0].Set(0, t)
|
||||
s.frame.Fields[1].Set(0, walker) // Value
|
||||
s.frame.Fields[2].Set(0, walker-((rand.Float64()*spread)+0.01)) // Min
|
||||
s.frame.Fields[3].Set(0, walker+((rand.Float64()*spread)+0.01)) // Max
|
||||
if err := sender.SendFrame(s.frame, mode); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user