Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a5c465a203 | ||
|
|
174b27004f | ||
|
|
a2e14ba6a9 | ||
|
|
a6683df10b | ||
|
|
2d4875fbd0 | ||
|
|
0a8f34423c |
@@ -0,0 +1 @@
|
|||||||
|
# Streaming Flow
|
||||||
+87
-93
@@ -1,103 +1,97 @@
|
|||||||
package cmd
|
package cmd
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"esdi/peripheral"
|
|
||||||
"fmt"
|
|
||||||
"strconv"
|
|
||||||
|
|
||||||
repl "github.com/ESilva15/ESgoRepl"
|
|
||||||
|
|
||||||
"github.com/spf13/cobra"
|
"github.com/spf13/cobra"
|
||||||
)
|
)
|
||||||
|
|
||||||
func replCmdAction(cmd *cobra.Command, args []string) {
|
func replCmdAction(cmd *cobra.Command, args []string) {
|
||||||
r := repl.NewREPL(repl.REPLCfg{
|
// r := repl.NewREPL(repl.REPLCfg{
|
||||||
PS1: "\rESDI > ",
|
// PS1: "\rESDI > ",
|
||||||
})
|
// })
|
||||||
|
//
|
||||||
perClerk := peripheral.NewPeripheralDeviceClerk()
|
// perClerk := peripheral.NewPeripheralDeviceClerk()
|
||||||
|
//
|
||||||
discoverDevicesREPLCmd := repl.Command{
|
// discoverDevicesREPLCmd := repl.Command{
|
||||||
Name: "discover",
|
// Name: "discover",
|
||||||
Usage: "discovers connected devices",
|
// Usage: "discovers connected devices",
|
||||||
Action: func(r *repl.REPL, args []string) error {
|
// Action: func(r *repl.REPL, args []string) error {
|
||||||
err := perClerk.FindDevices()
|
// err := perClerk.FindDevices()
|
||||||
if err != nil {
|
// if err != nil {
|
||||||
return err
|
// return err
|
||||||
}
|
// }
|
||||||
|
//
|
||||||
return nil
|
// return nil
|
||||||
},
|
// },
|
||||||
}
|
// }
|
||||||
|
//
|
||||||
listDevicesREPLCmd := repl.Command{
|
// listDevicesREPLCmd := repl.Command{
|
||||||
Name: "list",
|
// Name: "list",
|
||||||
Usage: "lists connected devices",
|
// Usage: "lists connected devices",
|
||||||
Action: func(r *repl.REPL, args []string) error {
|
// Action: func(r *repl.REPL, args []string) error {
|
||||||
_ = perClerk.ListDevices()
|
// _ = perClerk.ListDevices()
|
||||||
|
//
|
||||||
return nil
|
// return nil
|
||||||
},
|
// },
|
||||||
}
|
// }
|
||||||
|
//
|
||||||
listDeviceAPIREPLCmd := repl.Command{
|
// listDeviceAPIREPLCmd := repl.Command{
|
||||||
Name: "v-api",
|
// Name: "v-api",
|
||||||
Usage: "shows API of a device - pass its ID",
|
// Usage: "shows API of a device - pass its ID",
|
||||||
Action: func(r *repl.REPL, args []string) error {
|
// Action: func(r *repl.REPL, args []string) error {
|
||||||
// We should add this to the REPL instead
|
// // We should add this to the REPL instead
|
||||||
if len(args) < 1 {
|
// if len(args) < 1 {
|
||||||
return fmt.Errorf("requires at least on argument")
|
// return fmt.Errorf("requires at least on argument")
|
||||||
}
|
// }
|
||||||
|
//
|
||||||
// First and only argument should be the ID of the device we want to use
|
// // First and only argument should be the ID of the device we want to use
|
||||||
targetID, err := strconv.ParseInt(args[0], 10, 0)
|
// targetID, err := strconv.ParseInt(args[0], 10, 0)
|
||||||
if err != nil {
|
// if err != nil {
|
||||||
return err
|
// return err
|
||||||
}
|
// }
|
||||||
|
//
|
||||||
err = perClerk.ListDeviceAPI(uint8(targetID))
|
// err = perClerk.ListDeviceAPI(uint8(targetID))
|
||||||
if err != nil {
|
// if err != nil {
|
||||||
fmt.Println("failed to view device API: ", err.Error())
|
// fmt.Println("failed to view device API: ", err.Error())
|
||||||
}
|
// }
|
||||||
|
//
|
||||||
return nil
|
// return nil
|
||||||
},
|
// },
|
||||||
}
|
// }
|
||||||
|
//
|
||||||
runDeviceAPIREPLCmd := repl.Command{
|
// runDeviceAPIREPLCmd := repl.Command{
|
||||||
Name: "v-run",
|
// Name: "v-run",
|
||||||
Usage: "runs a funcion of a device - pass its ID and function name",
|
// Usage: "runs a funcion of a device - pass its ID and function name",
|
||||||
Action: func(r *repl.REPL, args []string) error {
|
// Action: func(r *repl.REPL, args []string) error {
|
||||||
// We should add this to the REPL instead
|
// // We should add this to the REPL instead
|
||||||
if len(args) < 3 {
|
// if len(args) < 3 {
|
||||||
return fmt.Errorf("requires at least on argument")
|
// return fmt.Errorf("requires at least on argument")
|
||||||
}
|
// }
|
||||||
|
//
|
||||||
// First and only argument should be the ID of the device we want to use
|
// // First and only argument should be the ID of the device we want to use
|
||||||
targetID, err := strconv.ParseInt(args[0], 10, 0)
|
// targetID, err := strconv.ParseInt(args[0], 10, 0)
|
||||||
if err != nil {
|
// if err != nil {
|
||||||
return err
|
// return err
|
||||||
}
|
// }
|
||||||
|
//
|
||||||
fnName := args[1]
|
// fnName := args[1]
|
||||||
fnArgs := args[2:]
|
// fnArgs := args[2:]
|
||||||
|
//
|
||||||
err = perClerk.RunDeviceFunction(uint8(targetID), fnName, fnArgs)
|
// err = perClerk.RunDeviceFunction(uint8(targetID), fnName, fnArgs)
|
||||||
if err != nil {
|
// if err != nil {
|
||||||
return err
|
// return err
|
||||||
}
|
// }
|
||||||
|
//
|
||||||
return nil
|
// return nil
|
||||||
},
|
// },
|
||||||
}
|
// }
|
||||||
|
//
|
||||||
r.RegisterCMD(discoverDevicesREPLCmd)
|
// r.RegisterCMD(discoverDevicesREPLCmd)
|
||||||
r.RegisterCMD(listDevicesREPLCmd)
|
// r.RegisterCMD(listDevicesREPLCmd)
|
||||||
r.RegisterCMD(listDeviceAPIREPLCmd)
|
// r.RegisterCMD(listDeviceAPIREPLCmd)
|
||||||
r.RegisterCMD(runDeviceAPIREPLCmd)
|
// r.RegisterCMD(runDeviceAPIREPLCmd)
|
||||||
|
//
|
||||||
r.Start()
|
// r.Start()
|
||||||
r.Close()
|
// r.Close()
|
||||||
}
|
}
|
||||||
|
|
||||||
// removeLabelCmd represents the removeLabel command
|
// removeLabelCmd represents the removeLabel command
|
||||||
|
|||||||
@@ -12,6 +12,7 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
helper "esdi/helpers"
|
helper "esdi/helpers"
|
||||||
|
"esdi/peripheral"
|
||||||
"esdi/peripheral/communication"
|
"esdi/peripheral/communication"
|
||||||
"esdi/peripheral/communication/packets"
|
"esdi/peripheral/communication/packets"
|
||||||
"esdi/peripheral/types"
|
"esdi/peripheral/types"
|
||||||
@@ -107,10 +108,12 @@ func NewCDashState() *CDashState {
|
|||||||
}
|
}
|
||||||
|
|
||||||
type CDashDisplay struct {
|
type CDashDisplay struct {
|
||||||
WT *communication.WalkieTalkie
|
WT *communication.WalkieTalkie
|
||||||
State *CDashState
|
State *CDashState
|
||||||
fieldToWindows map[telemetry.FieldID][]int16
|
fieldToWindows map[telemetry.FieldID][]int16
|
||||||
bufPool sync.Pool
|
bufPool sync.Pool
|
||||||
|
failedSends int
|
||||||
|
FailedSendsConsecutiveLimit int
|
||||||
}
|
}
|
||||||
|
|
||||||
// Connect will try to find and connect to the CDashDisplay
|
// Connect will try to find and connect to the CDashDisplay
|
||||||
@@ -118,7 +121,7 @@ func NewCDashDisplay() (*CDashDisplay, error) {
|
|||||||
// Look for the port
|
// Look for the port
|
||||||
p, err := findDisplayPort()
|
p, err := findDisplayPort()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
slog.Info("failed to find cdashdisplay port: %s", err.Error())
|
slog.Info("failed to find cdashdisplay port", "reason", err.Error())
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -132,9 +135,19 @@ func NewCDashDisplay() (*CDashDisplay, error) {
|
|||||||
return &b
|
return &b
|
||||||
},
|
},
|
||||||
},
|
},
|
||||||
|
failedSends: 0,
|
||||||
|
FailedSendsConsecutiveLimit: 5,
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (cds *CDashDisplay) Close() error {
|
||||||
|
// if cds.WT != nil {
|
||||||
|
// cds.Close()
|
||||||
|
// }
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
func (d *CDashDisplay) SendCommand() {
|
func (d *CDashDisplay) SendCommand() {
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -393,12 +406,12 @@ func (d *CDashDisplay) UnloadLayout() error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) {
|
func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) error {
|
||||||
packet := d.encodePacket(data)
|
packet := d.encodePacket(data)
|
||||||
|
|
||||||
bytes, err := helper.StructToBytes(packet)
|
bytes, err := helper.StructToBytes(packet)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return
|
return peripheral.ErrFailureToPackData
|
||||||
}
|
}
|
||||||
|
|
||||||
curStr := ""
|
curStr := ""
|
||||||
@@ -408,7 +421,6 @@ func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) {
|
|||||||
curStr += fmt.Sprintf("%02x ", byte)
|
curStr += fmt.Sprintf("%02x ", byte)
|
||||||
|
|
||||||
if byteCount == 8 {
|
if byteCount == 8 {
|
||||||
// slog.Debug(curStr)
|
|
||||||
curStr = ""
|
curStr = ""
|
||||||
byteCount = 0
|
byteCount = 0
|
||||||
}
|
}
|
||||||
@@ -417,6 +429,11 @@ func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) {
|
|||||||
// var ack packets.AckPacket
|
// var ack packets.AckPacket
|
||||||
err = d.WT.SendCommand(sendDataCMDID, bytes, nil)
|
err = d.WT.SendCommand(sendDataCMDID, bytes, nil)
|
||||||
if err != nil && err != io.EOF {
|
if err != nil && err != io.EOF {
|
||||||
return
|
if d.failedSends == d.FailedSendsConsecutiveLimit {
|
||||||
|
return peripheral.ErrDeviceTimedOut
|
||||||
|
}
|
||||||
|
d.failedSends++
|
||||||
}
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
}
|
}
|
||||||
|
|||||||
+1
-1
@@ -12,7 +12,7 @@ type Device struct {
|
|||||||
Discover func() (peripheral.Peripheral, error)
|
Discover func() (peripheral.Peripheral, error)
|
||||||
}
|
}
|
||||||
|
|
||||||
var List map[string]Device = map[string]Device{
|
var List map[string]*Device = map[string]*Device{
|
||||||
uidevice.NAME: {
|
uidevice.NAME: {
|
||||||
Name: uidevice.NAME,
|
Name: uidevice.NAME,
|
||||||
Discover: DiscoverUIDevice,
|
Discover: DiscoverUIDevice,
|
||||||
|
|||||||
@@ -17,9 +17,17 @@ func NewUIDevice() (peripheral.Peripheral, error) {
|
|||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (uid *UIDevice) SendData(data *telemetry.TelemetryData) {
|
func (uid *UIDevice) Close() error {
|
||||||
|
if uid.dataChan != nil {
|
||||||
|
close(uid.dataChan)
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (uid *UIDevice) SendData(data *telemetry.TelemetryData) error {
|
||||||
if data == nil {
|
if data == nil {
|
||||||
return
|
return peripheral.ErrInvalidData
|
||||||
}
|
}
|
||||||
|
|
||||||
select {
|
select {
|
||||||
@@ -27,6 +35,8 @@ func (uid *UIDevice) SendData(data *telemetry.TelemetryData) {
|
|||||||
default:
|
default:
|
||||||
// Drop frame if buffer is full
|
// Drop frame if buffer is full
|
||||||
}
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (uid *UIDevice) Name() string {
|
func (uid *UIDevice) Name() string {
|
||||||
|
|||||||
@@ -108,6 +108,15 @@ func (wt *WalkieTalkie) sendPacket(cmd types.Command, data any) error {
|
|||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// slog.Debug("# START ########################################################")
|
||||||
|
// slog.Debug(fmt.Sprintf("StartMarker: %02x", constvar.StartOfText))
|
||||||
|
// slog.Debug(fmt.Sprintf("CMD: %02x", cmd))
|
||||||
|
// slog.Debug(fmt.Sprintf("Len: %d", len(payload)))
|
||||||
|
// slog.Debug(fmt.Sprintf("Payload: %v", payload))
|
||||||
|
// slog.Debug(fmt.Sprintf("CRC: %v", CRC8(payload)))
|
||||||
|
// slog.Debug(fmt.Sprintf("EndMarker: %02x", constvar.EndOfText))
|
||||||
|
// slog.Debug("-")
|
||||||
|
|
||||||
packet := CMDDataPacket{
|
packet := CMDDataPacket{
|
||||||
StartMarker: constvar.StartOfText,
|
StartMarker: constvar.StartOfText,
|
||||||
CMD: cmd,
|
CMD: cmd,
|
||||||
@@ -119,7 +128,9 @@ func (wt *WalkieTalkie) sendPacket(cmd types.Command, data any) error {
|
|||||||
|
|
||||||
// Send the payload
|
// Send the payload
|
||||||
serializedPacket := packet.Serialize()
|
serializedPacket := packet.Serialize()
|
||||||
// fmt.Fprintf(os.Stderr, "%+v", serializedPacket)
|
|
||||||
|
// slog.Debug("Serialized packet", "packet", packet)
|
||||||
|
// slog.Debug("# END ##########################################################")
|
||||||
|
|
||||||
_, err = wt.Serial.Write(serializedPacket)
|
_, err = wt.Serial.Write(serializedPacket)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -176,21 +187,11 @@ func (wt *WalkieTalkie) readPacket(resp packets.Packet) error {
|
|||||||
// return nil
|
// return nil
|
||||||
// }
|
// }
|
||||||
|
|
||||||
func (wt *WalkieTalkie) SendCommand(cmd types.Command, payload any,
|
func (wt *WalkieTalkie) SendCommand(
|
||||||
responseBody packets.Packet) error {
|
cmd types.Command,
|
||||||
// Prepare the header
|
payload any,
|
||||||
// header := header{
|
responseBody packets.Packet,
|
||||||
// StartByte: constvar.StartOfText,
|
) error {
|
||||||
// CMD: cmd,
|
|
||||||
// EndByte: constvar.EndOfText,
|
|
||||||
// }
|
|
||||||
|
|
||||||
// Send the header
|
|
||||||
// err := wt.sendHeader(&header)
|
|
||||||
// if err != nil {
|
|
||||||
// return err
|
|
||||||
// }
|
|
||||||
|
|
||||||
// Send the body
|
// Send the body
|
||||||
err := wt.sendPacket(cmd, payload)
|
err := wt.sendPacket(cmd, payload)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -0,0 +1,9 @@
|
|||||||
|
package peripheral
|
||||||
|
|
||||||
|
import "errors"
|
||||||
|
|
||||||
|
var (
|
||||||
|
ErrInvalidData = errors.New("invalid data")
|
||||||
|
ErrDeviceTimedOut = errors.New("device timed out")
|
||||||
|
ErrFailureToPackData = errors.New("failed to pack received data")
|
||||||
|
)
|
||||||
@@ -17,8 +17,9 @@ const (
|
|||||||
|
|
||||||
type Peripheral interface {
|
type Peripheral interface {
|
||||||
Name() string
|
Name() string
|
||||||
SendData(*telemetry.TelemetryData)
|
SendData(*telemetry.TelemetryData) error
|
||||||
RequiredFields() []telemetry.FieldID
|
RequiredFields() []telemetry.FieldID
|
||||||
|
Close() error
|
||||||
}
|
}
|
||||||
|
|
||||||
type PeripheralDeviceClerk struct {
|
type PeripheralDeviceClerk struct {
|
||||||
|
|||||||
+17
-14
@@ -28,7 +28,7 @@ type BeamNG struct {
|
|||||||
updaters [telemetry.MaxFields]func(*telemetry.TelemetryField)
|
updaters [telemetry.MaxFields]func(*telemetry.TelemetryField)
|
||||||
|
|
||||||
// stream control
|
// stream control
|
||||||
streamCh chan telemetry.TelemetryData
|
wg sync.WaitGroup
|
||||||
streamCancel context.CancelFunc
|
streamCancel context.CancelFunc
|
||||||
|
|
||||||
// timing
|
// timing
|
||||||
@@ -46,12 +46,11 @@ func NewBeamNGProvider(logger *slog.Logger, opts *bngsdk.Options) (*BeamNG, erro
|
|||||||
}
|
}
|
||||||
|
|
||||||
provider := &BeamNG{
|
provider := &BeamNG{
|
||||||
logger: logger.With("TelemetryProvider", NAME),
|
logger: logger.With("TelemetryProvider", NAME),
|
||||||
streamCh: make(chan telemetry.TelemetryData, 1),
|
data: telemetry.NewTelemetryData(),
|
||||||
data: telemetry.NewTelemetryData(),
|
SDK: beam,
|
||||||
SDK: beam,
|
og: &bngsdk.Outgauge{},
|
||||||
og: &bngsdk.Outgauge{},
|
ticker: time.NewTicker(time.Second / 60),
|
||||||
ticker: time.NewTicker(time.Second / 60),
|
|
||||||
}
|
}
|
||||||
|
|
||||||
provider.updaters = [telemetry.MaxFields]func(*telemetry.TelemetryField){
|
provider.updaters = [telemetry.MaxFields]func(*telemetry.TelemetryField){
|
||||||
@@ -110,9 +109,9 @@ func (b *BeamNG) Stream() (<-chan telemetry.TelemetryData, error) {
|
|||||||
ctx, b.streamCancel = context.WithCancel(context.Background())
|
ctx, b.streamCancel = context.WithCancel(context.Background())
|
||||||
|
|
||||||
// Start the stream
|
// Start the stream
|
||||||
b.stream(ctx)
|
ch := b.stream(ctx)
|
||||||
|
|
||||||
return b.streamCh, nil
|
return ch, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *BeamNG) Subscribe(requestFields []telemetry.FieldID) {
|
func (b *BeamNG) Subscribe(requestFields []telemetry.FieldID) {
|
||||||
@@ -162,7 +161,6 @@ func (b *BeamNG) Subscribe(requestFields []telemetry.FieldID) {
|
|||||||
|
|
||||||
func (b *BeamNG) readData() {
|
func (b *BeamNG) readData() {
|
||||||
slog.Debug("READING THIS DATA")
|
slog.Debug("READING THIS DATA")
|
||||||
// BUG: getting stuck in here
|
|
||||||
ogSnapshot, err := b.SDK.Update()
|
ogSnapshot, err := b.SDK.Update()
|
||||||
slog.Debug("THE DATA WAS READ")
|
slog.Debug("THE DATA WAS READ")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -193,10 +191,15 @@ func (b *BeamNG) readData() {
|
|||||||
b.data.LastDataPoll = time.Now()
|
b.data.LastDataPoll = time.Now()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *BeamNG) stream(ctx context.Context) {
|
func (b *BeamNG) stream(ctx context.Context) <-chan telemetry.TelemetryData {
|
||||||
b.data.InitialTime = time.Now()
|
b.data.InitialTime = time.Now()
|
||||||
|
outCh := make(chan telemetry.TelemetryData)
|
||||||
|
b.wg.Add(1)
|
||||||
|
|
||||||
go func() {
|
go func() {
|
||||||
|
defer b.wg.Done()
|
||||||
|
defer close(outCh)
|
||||||
|
|
||||||
for {
|
for {
|
||||||
// Explicitly intercept cancellation
|
// Explicitly intercept cancellation
|
||||||
select {
|
select {
|
||||||
@@ -205,8 +208,6 @@ func (b *BeamNG) stream(ctx context.Context) {
|
|||||||
default:
|
default:
|
||||||
}
|
}
|
||||||
|
|
||||||
// NOTE: add a method to check if there's data available, or make this happen
|
|
||||||
|
|
||||||
select {
|
select {
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
return
|
return
|
||||||
@@ -217,7 +218,7 @@ func (b *BeamNG) stream(ctx context.Context) {
|
|||||||
|
|
||||||
// Publish data
|
// Publish data
|
||||||
select {
|
select {
|
||||||
case b.streamCh <- *b.data:
|
case outCh <- *b.data:
|
||||||
slog.Debug("PUBLISHED DATA")
|
slog.Debug("PUBLISHED DATA")
|
||||||
default:
|
default:
|
||||||
// skip this data, don't allow publishers to lag behind
|
// skip this data, don't allow publishers to lag behind
|
||||||
@@ -225,4 +226,6 @@ func (b *BeamNG) stream(ctx context.Context) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
|
return outCh
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -33,7 +33,8 @@ type IRacing struct {
|
|||||||
ticker *time.Ticker // ticker will keep polling intervals constant
|
ticker *time.Ticker // ticker will keep polling intervals constant
|
||||||
|
|
||||||
// Stream
|
// Stream
|
||||||
streamCh chan telemetry.TelemetryData
|
wg sync.WaitGroup
|
||||||
|
// streamCh chan telemetry.TelemetryData
|
||||||
streamCancel context.CancelFunc
|
streamCancel context.CancelFunc
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -50,13 +51,11 @@ func NewIRacingProvider(
|
|||||||
}
|
}
|
||||||
|
|
||||||
provider := &IRacing{
|
provider := &IRacing{
|
||||||
logger: logger,
|
logger: logger,
|
||||||
SDK: sdk,
|
SDK: sdk,
|
||||||
data: telemetry.NewTelemetryData(),
|
data: telemetry.NewTelemetryData(),
|
||||||
streamCh: make(chan telemetry.TelemetryData, 1),
|
|
||||||
// NOTE: This is because I stupidly recorded a test IBT file in 240
|
|
||||||
// TODO: make this configurable from the user side
|
// TODO: make this configurable from the user side
|
||||||
ticker: time.NewTicker(time.Second / 240),
|
ticker: time.NewTicker(time.Second / 60),
|
||||||
}
|
}
|
||||||
|
|
||||||
provider.updaters = [telemetry.MaxFields]func(*telemetry.TelemetryField){
|
provider.updaters = [telemetry.MaxFields]func(*telemetry.TelemetryField){
|
||||||
@@ -129,11 +128,14 @@ func (i *IRacing) isDataAvailable() bool {
|
|||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
|
|
||||||
func (i *IRacing) stream(ctx context.Context) {
|
func (i *IRacing) stream(ctx context.Context) <-chan telemetry.TelemetryData {
|
||||||
i.data.InitialTime = time.Now()
|
i.data.InitialTime = time.Now()
|
||||||
|
outCh := make(chan telemetry.TelemetryData)
|
||||||
|
i.wg.Add(1)
|
||||||
|
|
||||||
go func() {
|
go func() {
|
||||||
defer close(i.streamCh)
|
defer i.wg.Done()
|
||||||
|
defer close(outCh)
|
||||||
|
|
||||||
// Put this into the configuration file
|
// Put this into the configuration file
|
||||||
consecutiveTimeouts := 0
|
consecutiveTimeouts := 0
|
||||||
@@ -149,12 +151,11 @@ func (i *IRacing) stream(ctx context.Context) {
|
|||||||
|
|
||||||
if i.SDK.CheckForDataEvent(time.Duration(dataEvTimeout) * time.Millisecond) {
|
if i.SDK.CheckForDataEvent(time.Duration(dataEvTimeout) * time.Millisecond) {
|
||||||
consecutiveTimeouts = 0
|
consecutiveTimeouts = 0
|
||||||
i.logger.Debug("sending data", "timeouts", consecutiveTimeouts)
|
|
||||||
i.readData()
|
i.readData()
|
||||||
|
|
||||||
// Publish data
|
// Publish data
|
||||||
select {
|
select {
|
||||||
case i.streamCh <- *i.data:
|
case outCh <- *i.data:
|
||||||
default:
|
default:
|
||||||
// skip this data, don't allow publishers to lag behind
|
// skip this data, don't allow publishers to lag behind
|
||||||
}
|
}
|
||||||
@@ -170,6 +171,8 @@ func (i *IRacing) stream(ctx context.Context) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
|
||||||
|
return outCh
|
||||||
}
|
}
|
||||||
|
|
||||||
func (i *IRacing) readData() {
|
func (i *IRacing) readData() {
|
||||||
@@ -206,9 +209,9 @@ func (i *IRacing) Stream() (<-chan telemetry.TelemetryData, error) {
|
|||||||
ctx, i.streamCancel = context.WithCancel(context.Background())
|
ctx, i.streamCancel = context.WithCancel(context.Background())
|
||||||
|
|
||||||
// Start the stream
|
// Start the stream
|
||||||
i.stream(ctx)
|
ch := i.stream(ctx)
|
||||||
|
|
||||||
return i.streamCh, nil
|
return ch, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (i *IRacing) StopStream() {
|
func (i *IRacing) StopStream() {
|
||||||
@@ -217,6 +220,7 @@ func (i *IRacing) StopStream() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
i.streamCancel()
|
i.streamCancel()
|
||||||
|
i.wg.Wait()
|
||||||
i.streamCancel = nil
|
i.streamCancel = nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+32
-78
@@ -2,32 +2,21 @@ package services
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
|
||||||
"fmt"
|
|
||||||
"log/slog"
|
"log/slog"
|
||||||
"sync"
|
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
"time"
|
|
||||||
|
|
||||||
"esdi/devices"
|
"esdi/devices"
|
||||||
"esdi/peripheral"
|
"esdi/peripheral"
|
||||||
"esdi/telemetry"
|
"esdi/telemetry"
|
||||||
)
|
)
|
||||||
|
|
||||||
var ErrPeripheralAlreadyRegistered = errors.New("peripheral is already registered")
|
// DeviceService is the API for the peripherals
|
||||||
|
|
||||||
// DeviceService will handle sending the data from the telemetry service to the
|
|
||||||
// actual devices
|
|
||||||
// NOTE: create a virtual device and make it be the output window or something so
|
|
||||||
// we can just add it as a device or whatever instead of being a custom made thing
|
|
||||||
// that would be pretty cool I think
|
|
||||||
type DeviceService struct {
|
type DeviceService struct {
|
||||||
Logger *slog.Logger
|
Logger *slog.Logger
|
||||||
// Device discovery
|
// Device discovery
|
||||||
mu sync.RWMutex
|
PSS *PeripheralStateStore // Store to track peripheral state
|
||||||
ctxDiscovery context.Context
|
ctxDiscovery context.Context
|
||||||
ctxDiscoveryCancel context.CancelFunc
|
ctxDiscoveryCancel context.CancelFunc
|
||||||
Devices map[string]peripheral.Peripheral
|
|
||||||
// Strem handling
|
// Strem handling
|
||||||
streamCancel context.CancelFunc
|
streamCancel context.CancelFunc
|
||||||
TelemCh <-chan telemetry.TelemetryData
|
TelemCh <-chan telemetry.TelemetryData
|
||||||
@@ -39,7 +28,7 @@ func NewDeviceService(logger *slog.Logger) *DeviceService {
|
|||||||
sharedChannel := make(chan string, 10)
|
sharedChannel := make(chan string, 10)
|
||||||
|
|
||||||
dev := &DeviceService{
|
dev := &DeviceService{
|
||||||
Devices: make(map[string]peripheral.Peripheral),
|
PSS: NewPeripheralStateStore(logger.With("Service", "PeripheralStateStore"), devices.List),
|
||||||
Logger: logger,
|
Logger: logger,
|
||||||
Messages: sharedChannel,
|
Messages: sharedChannel,
|
||||||
}
|
}
|
||||||
@@ -52,76 +41,33 @@ func NewDeviceService(logger *slog.Logger) *DeviceService {
|
|||||||
return dev
|
return dev
|
||||||
}
|
}
|
||||||
|
|
||||||
func (ds *DeviceService) FindDevices() {
|
// Getters [START] -------------------------------------------------------------
|
||||||
// Need to define a list of devices to search for
|
func (ds *DeviceService) GetDevices() []peripheral.Peripheral {
|
||||||
// For now lets just try to find our cdashdisplay - will think about the rest later
|
snapshot := ds.PSS.GetStates()
|
||||||
ticker := time.NewTicker(2 * time.Second)
|
peripherals := make([]peripheral.Peripheral, 0, len(snapshot))
|
||||||
defer ticker.Stop()
|
|
||||||
|
|
||||||
for {
|
for _, state := range snapshot {
|
||||||
select {
|
if state.State != DeviceIsConnected {
|
||||||
case <-ds.ctxDiscovery.Done():
|
continue
|
||||||
// If requested to cancel we cancel background discovery
|
|
||||||
return
|
|
||||||
case <-ticker.C:
|
|
||||||
for pName, peripheral := range devices.List {
|
|
||||||
if ds.DeviceExists(pName) {
|
|
||||||
// We already discovered this device
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
ds.Logger.Debug("looking for device", "name", pName)
|
|
||||||
dev, err := peripheral.Discover()
|
|
||||||
if err != nil {
|
|
||||||
ds.Logger.Debug("didn't find device", "name", pName)
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
|
|
||||||
// Register the device we just found
|
|
||||||
ds.RegisterDevice(dev)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
peripherals = append(peripherals, state.Peripheral)
|
||||||
}
|
|
||||||
|
|
||||||
// func (ds *DeviceService) SubscribeFields() error {
|
|
||||||
// for _, dev := range ds.Devices {
|
|
||||||
// fields := dev.RequiredFields()
|
|
||||||
// }
|
|
||||||
//
|
|
||||||
// return nil
|
|
||||||
// }
|
|
||||||
|
|
||||||
func (ds *DeviceService) RegisterDevice(dev peripheral.Peripheral) error {
|
|
||||||
ds.mu.Lock()
|
|
||||||
defer ds.mu.Unlock()
|
|
||||||
|
|
||||||
if ds.DeviceExists(dev.Name()) {
|
|
||||||
return ErrPeripheralAlreadyRegistered
|
|
||||||
}
|
}
|
||||||
|
|
||||||
ds.Devices[dev.Name()] = dev
|
return peripherals
|
||||||
|
|
||||||
return nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (ds *DeviceService) GetDevice(name string) (peripheral.Peripheral, error) {
|
func (ds *DeviceService) GetPeripheral(pname string) (peripheral.Peripheral, error) {
|
||||||
val, ok := ds.Devices[name]
|
return ds.PSS.GetPeripheral(pname)
|
||||||
if !ok {
|
|
||||||
return nil, fmt.Errorf("device `%s` couldn't be found", name)
|
|
||||||
}
|
|
||||||
|
|
||||||
return val, nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (ds *DeviceService) DeviceExists(name string) bool {
|
func (ds *DeviceService) PeripheralExists(pname string) bool {
|
||||||
if _, ok := ds.Devices[name]; !ok {
|
_, err := ds.PSS.GetPeripheral(pname)
|
||||||
return false
|
return err == nil
|
||||||
}
|
|
||||||
|
|
||||||
return true
|
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Getters [END] ---------------------------------------------------------------
|
||||||
|
|
||||||
|
// Actions [START] -------------------------------------------------------------
|
||||||
func (ds *DeviceService) StartStream() {
|
func (ds *DeviceService) StartStream() {
|
||||||
// NOTE: i'm using this pattern a whole lot. Maybe I can create a struct to handle this
|
// NOTE: i'm using this pattern a whole lot. Maybe I can create a struct to handle this
|
||||||
var ctx context.Context
|
var ctx context.Context
|
||||||
@@ -144,6 +90,8 @@ func (ds *DeviceService) SetTelemetryChannel(ch <-chan telemetry.TelemetryData)
|
|||||||
ds.TelemCh = ch
|
ds.TelemCh = ch
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Actions [END] ---------------------------------------------------------------
|
||||||
|
|
||||||
// transmit will send the data to the devices themselves
|
// transmit will send the data to the devices themselves
|
||||||
func (ds *DeviceService) transmit(ctx context.Context) {
|
func (ds *DeviceService) transmit(ctx context.Context) {
|
||||||
var isSending atomic.Bool
|
var isSending atomic.Bool
|
||||||
@@ -165,11 +113,17 @@ func (ds *DeviceService) transmit(ctx context.Context) {
|
|||||||
|
|
||||||
// TODO: make a copy of the data and send that copy instead of keeping
|
// TODO: make a copy of the data and send that copy instead of keeping
|
||||||
// the data locked
|
// the data locked
|
||||||
ds.mu.RLock()
|
for _, dev := range ds.PSS.GetStates() {
|
||||||
for _, dev := range ds.Devices {
|
if dev.State != DeviceIsConnected {
|
||||||
dev.SendData(&data)
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
err := dev.Peripheral.SendData(&data)
|
||||||
|
|
||||||
|
if err == peripheral.ErrDeviceTimedOut {
|
||||||
|
ds.onDeviceTimedOut(dev.device.Name)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
ds.mu.RUnlock()
|
|
||||||
|
|
||||||
isSending.Store(false)
|
isSending.Store(false)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,222 @@
|
|||||||
|
package services
|
||||||
|
|
||||||
|
import (
|
||||||
|
"errors"
|
||||||
|
"log/slog"
|
||||||
|
"maps"
|
||||||
|
"sync"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"esdi/devices"
|
||||||
|
"esdi/peripheral"
|
||||||
|
)
|
||||||
|
|
||||||
|
var (
|
||||||
|
ErrPeripheralAlreadyRegistered = errors.New("peripheral is already registered")
|
||||||
|
ErrNoSuchDevice = errors.New("device doesn't exist")
|
||||||
|
ErrDeviceIsNotConnected = errors.New("device isn't connected")
|
||||||
|
)
|
||||||
|
|
||||||
|
type DeviceState = uint8
|
||||||
|
|
||||||
|
const (
|
||||||
|
DeviceTimedOut uint8 = iota
|
||||||
|
DeviceIsConnected
|
||||||
|
DeviceIsDisconnected
|
||||||
|
DeviceReconnected
|
||||||
|
)
|
||||||
|
|
||||||
|
type PeripheralState struct {
|
||||||
|
device *devices.Device
|
||||||
|
Peripheral peripheral.Peripheral
|
||||||
|
State DeviceState
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewPeripheralState(
|
||||||
|
dev *devices.Device,
|
||||||
|
peripheral peripheral.Peripheral,
|
||||||
|
state DeviceState,
|
||||||
|
) *PeripheralState {
|
||||||
|
perState := PeripheralState{
|
||||||
|
device: dev,
|
||||||
|
Peripheral: peripheral,
|
||||||
|
State: state,
|
||||||
|
}
|
||||||
|
|
||||||
|
return &perState
|
||||||
|
}
|
||||||
|
|
||||||
|
type PeripheralStateStore struct {
|
||||||
|
Logger *slog.Logger
|
||||||
|
mu sync.RWMutex
|
||||||
|
store map[string]*PeripheralState
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewPeripheralStateStore(
|
||||||
|
nLogger *slog.Logger,
|
||||||
|
devList map[string]*devices.Device,
|
||||||
|
) *PeripheralStateStore {
|
||||||
|
store := PeripheralStateStore{
|
||||||
|
Logger: nLogger,
|
||||||
|
store: make(map[string]*PeripheralState),
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, dev := range devList {
|
||||||
|
store.AddDevice(dev)
|
||||||
|
}
|
||||||
|
|
||||||
|
return &store
|
||||||
|
}
|
||||||
|
|
||||||
|
func (pss *PeripheralStateStore) GetStates() map[string]*PeripheralState {
|
||||||
|
pss.mu.RLock()
|
||||||
|
defer pss.mu.RUnlock()
|
||||||
|
|
||||||
|
return maps.Clone(pss.store)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (pss *PeripheralStateStore) GetState(pname string) (*PeripheralState, error) {
|
||||||
|
if !pss.DeviceExists(pname) {
|
||||||
|
return nil, ErrNoSuchDevice
|
||||||
|
}
|
||||||
|
|
||||||
|
pss.mu.RLock()
|
||||||
|
defer pss.mu.RUnlock()
|
||||||
|
return pss.store[pname], nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (pss *PeripheralStateStore) GetPeripheral(pname string) (peripheral.Peripheral, error) {
|
||||||
|
state, err := pss.GetState(pname)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
|
||||||
|
if state.State != DeviceIsConnected {
|
||||||
|
return nil, ErrDeviceIsNotConnected
|
||||||
|
}
|
||||||
|
|
||||||
|
return state.Peripheral, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// AddDevice adds a new device for tracking
|
||||||
|
func (pss *PeripheralStateStore) AddDevice(dev *devices.Device) error {
|
||||||
|
if pss.DeviceExists(dev.Name) {
|
||||||
|
return ErrPeripheralAlreadyRegistered
|
||||||
|
}
|
||||||
|
|
||||||
|
pss.mu.Lock()
|
||||||
|
defer pss.mu.Unlock()
|
||||||
|
pss.store[dev.Name] = NewPeripheralState(dev, nil, DeviceIsDisconnected)
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// DeviceExists returns whether the store is already tracking `pname`
|
||||||
|
func (pss *PeripheralStateStore) DeviceExists(pname string) bool {
|
||||||
|
pss.mu.RLock()
|
||||||
|
defer pss.mu.RUnlock()
|
||||||
|
|
||||||
|
if _, ok := pss.store[pname]; ok {
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
// DeleteDevice deletes `pname` from tracking
|
||||||
|
func (pss *PeripheralStateStore) DeleteDevice(pname string) error {
|
||||||
|
if !pss.DeviceExists(pname) {
|
||||||
|
return ErrNoSuchDevice
|
||||||
|
}
|
||||||
|
|
||||||
|
pss.mu.Lock()
|
||||||
|
defer pss.mu.Unlock()
|
||||||
|
|
||||||
|
delete(pss.store, pname)
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// "Events" [START] ------------------------------------------------------------
|
||||||
|
func (ds *DeviceService) onDeviceTimedOut(pname string) {
|
||||||
|
// We need to deregister the device
|
||||||
|
ds.PSS.setDeviceTimedOut(pname)
|
||||||
|
}
|
||||||
|
|
||||||
|
// "Events" [END] --------------------------------------------------------------
|
||||||
|
|
||||||
|
// Device State Handling [START] -----------------------------------------------
|
||||||
|
func (pss *PeripheralStateStore) setDeviceConnected(pname string, per peripheral.Peripheral) {
|
||||||
|
pss.Logger.Info("found device", "device", pname)
|
||||||
|
|
||||||
|
pss.mu.Lock()
|
||||||
|
defer pss.mu.Unlock()
|
||||||
|
pss.store[pname].Peripheral = per
|
||||||
|
pss.store[pname].State = DeviceIsConnected
|
||||||
|
}
|
||||||
|
|
||||||
|
func (pss *PeripheralStateStore) setDeviceTimedOut(pname string) {
|
||||||
|
pss.Logger.Info("device timed out", "device", pname)
|
||||||
|
|
||||||
|
pss.mu.Lock()
|
||||||
|
defer pss.mu.Unlock()
|
||||||
|
pss.store[pname].Peripheral = nil
|
||||||
|
pss.store[pname].State = DeviceTimedOut
|
||||||
|
}
|
||||||
|
|
||||||
|
func (pss *PeripheralStateStore) setDeviceReconnected(pname string, per peripheral.Peripheral) {
|
||||||
|
pss.Logger.Info("device reconnecting", "device", pname)
|
||||||
|
|
||||||
|
pss.mu.Lock()
|
||||||
|
defer pss.mu.Unlock()
|
||||||
|
pss.store[pname].Peripheral = per
|
||||||
|
pss.store[pname].State = DeviceReconnected
|
||||||
|
}
|
||||||
|
|
||||||
|
// Device State Handling [END] -------------------------------------------------
|
||||||
|
|
||||||
|
// Device Handling [START] -----------------------------------------------------
|
||||||
|
func (pss *PeripheralStateStore) HandleDeviceState() {
|
||||||
|
peripherals := pss.GetStates()
|
||||||
|
|
||||||
|
for pName, pState := range peripherals {
|
||||||
|
switch pState.State {
|
||||||
|
case DeviceIsDisconnected:
|
||||||
|
pss.Logger.Debug("looking for device", "name", pName)
|
||||||
|
dev, err := pState.device.Discover()
|
||||||
|
if err != nil {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
// Register the device we just found
|
||||||
|
pss.setDeviceConnected(pName, dev)
|
||||||
|
case DeviceIsConnected:
|
||||||
|
// Need to check if its streaming, if its not streaming than we have to do a healthcheck
|
||||||
|
pss.Logger.Debug("Device is connected. Normal", "device", pName)
|
||||||
|
case DeviceReconnected:
|
||||||
|
// If the device has reconnected we need to reset the device and then set it as connected
|
||||||
|
pss.Logger.Debug("Device has reconnected. Clearing up state", "device", pName)
|
||||||
|
case DeviceTimedOut:
|
||||||
|
// If the device has timed out we need to re-discover it or something
|
||||||
|
pss.Logger.Debug("Device is timed out. Attempting to recconect", "device", pName)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Device Handling [END] -------------------------------------------------------
|
||||||
|
|
||||||
|
// FindDevices is a routine that goes over the devices in the PeripheralStateStore
|
||||||
|
// and handles their state accordingly
|
||||||
|
func (ds *DeviceService) FindDevices() {
|
||||||
|
ticker := time.NewTicker(2 * time.Second)
|
||||||
|
defer ticker.Stop()
|
||||||
|
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-ds.ctxDiscovery.Done():
|
||||||
|
return
|
||||||
|
case <-ticker.C:
|
||||||
|
ds.PSS.HandleDeviceState()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
+125
-103
@@ -16,6 +16,8 @@ import (
|
|||||||
type TelemetryService struct {
|
type TelemetryService struct {
|
||||||
logger *slog.Logger
|
logger *slog.Logger
|
||||||
devService *DeviceService
|
devService *DeviceService
|
||||||
|
// Streaming
|
||||||
|
isStreaming bool
|
||||||
// Concurrency protection
|
// Concurrency protection
|
||||||
mut sync.RWMutex
|
mut sync.RWMutex
|
||||||
activeProvider telem.TelemetryProvider
|
activeProvider telem.TelemetryProvider
|
||||||
@@ -66,6 +68,92 @@ func (t *TelemetryService) ProviderMonitor(ctx context.Context) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Listener Control [START] ----------------------------------------------------
|
||||||
|
|
||||||
|
func (t *TelemetryService) SubscribeListener(id string, bufferSize int) <-chan telem.TelemetryData {
|
||||||
|
t.mut.Lock()
|
||||||
|
defer t.mut.Unlock()
|
||||||
|
|
||||||
|
// NOTE: is this truly necessary?
|
||||||
|
// return the channel if it already exists
|
||||||
|
if ch, exists := t.listeners[id]; exists {
|
||||||
|
return ch
|
||||||
|
}
|
||||||
|
|
||||||
|
ch := make(chan telem.TelemetryData, bufferSize)
|
||||||
|
t.listeners[id] = ch
|
||||||
|
|
||||||
|
t.logger.Info("New stream subscriber registered", "id", id)
|
||||||
|
return ch
|
||||||
|
}
|
||||||
|
|
||||||
|
func (t *TelemetryService) UnsubscribeListener(id string) {
|
||||||
|
t.mut.Lock()
|
||||||
|
defer t.mut.Unlock()
|
||||||
|
|
||||||
|
if ch, exists := t.listeners[id]; exists {
|
||||||
|
close(ch)
|
||||||
|
delete(t.listeners, id)
|
||||||
|
t.logger.Info("Stream subscriber removed", "id", id)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (t *TelemetryService) SubscribeToFields() []telem.FieldID {
|
||||||
|
seen := make(map[telemetry.FieldID]struct{})
|
||||||
|
var allFields []telemetry.FieldID
|
||||||
|
|
||||||
|
for _, dev := range t.devService.GetDevices() {
|
||||||
|
for _, field := range dev.RequiredFields() {
|
||||||
|
if _, exists := seen[field]; !exists {
|
||||||
|
seen[field] = struct{}{}
|
||||||
|
allFields = append(allFields, field)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
t.logger.Debug("requested fields", "fields", allFields)
|
||||||
|
t.activeProvider.Subscribe(allFields)
|
||||||
|
|
||||||
|
return allFields
|
||||||
|
}
|
||||||
|
|
||||||
|
// Listener Control [END] ------------------------------------------------------
|
||||||
|
|
||||||
|
// Provider Control [START] ----------------------------------------------------
|
||||||
|
|
||||||
|
func (t *TelemetryService) HasActiveProvider() bool {
|
||||||
|
if t.activeProvider == nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
func (t *TelemetryService) dropActiveProvider() {
|
||||||
|
if t.cancelForward != nil {
|
||||||
|
t.cancelForward()
|
||||||
|
}
|
||||||
|
|
||||||
|
t.activeProvider.StopStream()
|
||||||
|
t.activeProvider.Close()
|
||||||
|
t.activeProvider = nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (t *TelemetryService) SwitchProvider(newProvider telem.TelemetryProvider) error {
|
||||||
|
t.mut.Lock()
|
||||||
|
defer t.mut.Unlock()
|
||||||
|
|
||||||
|
// Clean up the current to be old provider
|
||||||
|
if t.activeProvider != nil {
|
||||||
|
t.dropActiveProvider()
|
||||||
|
}
|
||||||
|
|
||||||
|
// Assign the new provider
|
||||||
|
t.activeProvider = newProvider
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
func (t *TelemetryService) onProviderHealthCheckFailed() {
|
func (t *TelemetryService) onProviderHealthCheckFailed() {
|
||||||
// Just restart the whole lookup process
|
// Just restart the whole lookup process
|
||||||
go t.FindProvider(t.CtxMonitor)
|
go t.FindProvider(t.CtxMonitor)
|
||||||
@@ -123,19 +211,46 @@ func (t *TelemetryService) FindProvider(ctx context.Context) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *TelemetryService) SwitchProvider(newProvider telem.TelemetryProvider) error {
|
// Provider Control [END] ------------------------------------------------------
|
||||||
|
|
||||||
|
// Streaming Control [START] ---------------------------------------------------
|
||||||
|
|
||||||
|
func (t *TelemetryService) StopStream() {
|
||||||
t.mut.Lock()
|
t.mut.Lock()
|
||||||
defer t.mut.Unlock()
|
defer t.mut.Unlock()
|
||||||
|
|
||||||
// Clean up the current to be old provider
|
if t.cancelForward != nil {
|
||||||
if t.activeProvider != nil {
|
t.cancelForward()
|
||||||
t.dropActiveProvider()
|
t.cancelForward = nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Assign the new provider
|
t.activeProvider.StopStream()
|
||||||
t.activeProvider = newProvider
|
t.isStreaming = false
|
||||||
|
}
|
||||||
|
|
||||||
return nil
|
func (t *TelemetryService) StartStream() {
|
||||||
|
slog.Debug("Stream started")
|
||||||
|
|
||||||
|
// Start the new stream
|
||||||
|
if t.activeProvider == nil {
|
||||||
|
slog.Debug("there's no active provider. not starting the stream")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
// Stop the provider healthcheck
|
||||||
|
t.healthCheckCancel()
|
||||||
|
|
||||||
|
simInCh, _ := t.activeProvider.Stream()
|
||||||
|
// TODO: the provider needs to be able to tell the data has stopped
|
||||||
|
// so we can restart the provider lookup routine
|
||||||
|
|
||||||
|
// Create the context so we can control the lifecycle
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
t.cancelForward = cancel
|
||||||
|
|
||||||
|
// Multiplex this data
|
||||||
|
go t.multiplexData(ctx, simInCh)
|
||||||
|
t.isStreaming = true
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *TelemetryService) multiplexData(ctx context.Context, dataCh <-chan telem.TelemetryData) {
|
func (t *TelemetryService) multiplexData(ctx context.Context, dataCh <-chan telem.TelemetryData) {
|
||||||
@@ -165,101 +280,8 @@ func (t *TelemetryService) multiplexData(ctx context.Context, dataCh <-chan tele
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *TelemetryService) dropActiveProvider() {
|
func (t *TelemetryService) IsStreaming() bool {
|
||||||
if t.cancelForward != nil {
|
return t.isStreaming
|
||||||
t.cancelForward()
|
|
||||||
}
|
|
||||||
|
|
||||||
t.activeProvider.StopStream()
|
|
||||||
t.activeProvider.Close()
|
|
||||||
t.activeProvider = nil
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (t *TelemetryService) SubscribeListener(id string, bufferSize int) <-chan telem.TelemetryData {
|
// Streaming Control [END] -----------------------------------------------------
|
||||||
t.mut.Lock()
|
|
||||||
defer t.mut.Unlock()
|
|
||||||
|
|
||||||
// NOTE: is this truly necessary?
|
|
||||||
// return the channel if it already exists
|
|
||||||
if ch, exists := t.listeners[id]; exists {
|
|
||||||
return ch
|
|
||||||
}
|
|
||||||
|
|
||||||
ch := make(chan telem.TelemetryData, bufferSize)
|
|
||||||
t.listeners[id] = ch
|
|
||||||
|
|
||||||
t.logger.Info("New stream subscriber registered", "id", id)
|
|
||||||
return ch
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t *TelemetryService) UnsubscribeListener(id string) {
|
|
||||||
t.mut.Lock()
|
|
||||||
defer t.mut.Unlock()
|
|
||||||
|
|
||||||
if ch, exists := t.listeners[id]; exists {
|
|
||||||
close(ch)
|
|
||||||
delete(t.listeners, id)
|
|
||||||
t.logger.Info("Stream subscriber removed", "id", id)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t *TelemetryService) SubscribeToFields() {
|
|
||||||
seen := make(map[telemetry.FieldID]struct{})
|
|
||||||
var allFields []telemetry.FieldID
|
|
||||||
|
|
||||||
for _, dev := range t.devService.Devices {
|
|
||||||
for _, field := range dev.RequiredFields() {
|
|
||||||
if _, exists := seen[field]; !exists {
|
|
||||||
seen[field] = struct{}{}
|
|
||||||
allFields = append(allFields, field)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
t.logger.Debug("requested fields", "fields", allFields)
|
|
||||||
t.activeProvider.Subscribe(allFields)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t *TelemetryService) StartStream() {
|
|
||||||
slog.Debug("Stream started")
|
|
||||||
|
|
||||||
// Start the new stream
|
|
||||||
if t.activeProvider == nil {
|
|
||||||
slog.Debug("there's no active provider. not starting the stream")
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
// Stop the provider healthcheck
|
|
||||||
t.healthCheckCancel()
|
|
||||||
|
|
||||||
simInCh, _ := t.activeProvider.Stream()
|
|
||||||
// TODO: the provider needs to be able to tell the data has stopped
|
|
||||||
// so we can restart the provider lookup routine
|
|
||||||
|
|
||||||
// Create the context so we can control the lifecycle
|
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
|
||||||
t.cancelForward = cancel
|
|
||||||
|
|
||||||
// Multiplex this data
|
|
||||||
go t.multiplexData(ctx, simInCh)
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t *TelemetryService) StopStream() {
|
|
||||||
t.mut.Lock()
|
|
||||||
defer t.mut.Unlock()
|
|
||||||
|
|
||||||
if t.cancelForward != nil {
|
|
||||||
t.cancelForward()
|
|
||||||
t.cancelForward = nil
|
|
||||||
}
|
|
||||||
|
|
||||||
t.activeProvider.StopStream()
|
|
||||||
}
|
|
||||||
|
|
||||||
func (t *TelemetryService) HasActiveProvider() bool {
|
|
||||||
if t.activeProvider == nil {
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
|
|
||||||
return true
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -211,7 +211,7 @@ func (lc *LayoutController) createWindow() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Acquire the cdashdisplay
|
// Acquire the cdashdisplay
|
||||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||||
return
|
return
|
||||||
@@ -321,7 +321,7 @@ func (lc *LayoutController) newWindowAction() {
|
|||||||
|
|
||||||
func (lc *LayoutController) updateWindowAction(win *cdashdisplay.DesktopUIWindow) {
|
func (lc *LayoutController) updateWindowAction(win *cdashdisplay.DesktopUIWindow) {
|
||||||
// Acquire the cdashdisplay
|
// Acquire the cdashdisplay
|
||||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||||
return
|
return
|
||||||
@@ -347,7 +347,7 @@ func (lc *LayoutController) displayLoadedLayouts() {
|
|||||||
lc.Logger.Debug("We want to view our layout!")
|
lc.Logger.Debug("We want to view our layout!")
|
||||||
|
|
||||||
// Acquire the cdashdisplay
|
// Acquire the cdashdisplay
|
||||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||||
return
|
return
|
||||||
@@ -392,7 +392,7 @@ func (lc *LayoutController) getCurrentTreeNodeModel() (*tview.TreeNode, int16, e
|
|||||||
|
|
||||||
func (lc *LayoutController) loadLayout() {
|
func (lc *LayoutController) loadLayout() {
|
||||||
// Acquire the cdashdisplay
|
// Acquire the cdashdisplay
|
||||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||||
return
|
return
|
||||||
@@ -416,7 +416,7 @@ func (lc *LayoutController) loadLayout() {
|
|||||||
|
|
||||||
func (lc *LayoutController) unloadLayout() {
|
func (lc *LayoutController) unloadLayout() {
|
||||||
// Acquire the cdashdisplay
|
// Acquire the cdashdisplay
|
||||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||||
return
|
return
|
||||||
@@ -437,7 +437,7 @@ func (lc *LayoutController) unloadLayout() {
|
|||||||
|
|
||||||
func (lc *LayoutController) saveLayout() {
|
func (lc *LayoutController) saveLayout() {
|
||||||
// Acquire the cdashdisplay
|
// Acquire the cdashdisplay
|
||||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||||
return
|
return
|
||||||
@@ -470,7 +470,7 @@ func (lc *LayoutController) deleteWindow() {
|
|||||||
wID := curNode.GetReference().(int16)
|
wID := curNode.GetReference().(int16)
|
||||||
|
|
||||||
// Acquire the cdashdisplay
|
// Acquire the cdashdisplay
|
||||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||||
return
|
return
|
||||||
|
|||||||
@@ -52,7 +52,7 @@ func (lc *LayoutController) handleMovementCapture(idx int16,
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Acquire the cdashdisplay
|
// Acquire the cdashdisplay
|
||||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||||
return nil
|
return nil
|
||||||
@@ -86,7 +86,7 @@ func (lc *LayoutController) handleResizeCapture(idx int16,
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Acquire the cdashdisplay
|
// Acquire the cdashdisplay
|
||||||
displayIF, err := lc.DevService.GetDevice(cdashdisplay.NAME)
|
displayIF, err := lc.DevService.GetPeripheral(cdashdisplay.NAME)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
lc.Messages <- "failed to get " + cdashdisplay.NAME
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
@@ -67,7 +67,7 @@ func (mc *DeviceController) AddDeviceAPIListItems() {
|
|||||||
mc.DeviceAPIView.DevAPIList.
|
mc.DeviceAPIView.DevAPIList.
|
||||||
AddItem("layout", "build a layout for CDashDisplay", func() {
|
AddItem("layout", "build a layout for CDashDisplay", func() {
|
||||||
// This CDashDisplay specific, only load if we have a CDashDisplay
|
// This CDashDisplay specific, only load if we have a CDashDisplay
|
||||||
if !mc.DevService.DeviceExists(cdashdisplay.NAME) {
|
if !mc.DevService.PeripheralExists(cdashdisplay.NAME) {
|
||||||
mc.DevService.Messages <- "CDashDisplay it not loaded yet\n"
|
mc.DevService.Messages <- "CDashDisplay it not loaded yet\n"
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -18,7 +18,7 @@ import (
|
|||||||
|
|
||||||
type StreamingCtrl struct {
|
type StreamingCtrl struct {
|
||||||
*Controller
|
*Controller
|
||||||
Service *services.DeviceService
|
DevService *services.DeviceService
|
||||||
StreamView *views.StreamToolView
|
StreamView *views.StreamToolView
|
||||||
Messages chan string
|
Messages chan string
|
||||||
Internal chan string
|
Internal chan string
|
||||||
@@ -45,18 +45,16 @@ func NewStreamingCtrl(
|
|||||||
|
|
||||||
ctrl := &StreamingCtrl{
|
ctrl := &StreamingCtrl{
|
||||||
Controller: base,
|
Controller: base,
|
||||||
Service: devService,
|
DevService: devService,
|
||||||
TelemServ: serTelem,
|
TelemServ: serTelem,
|
||||||
Messages: make(chan string, 10),
|
Messages: make(chan string, 10),
|
||||||
Internal: make(chan string, 10),
|
Internal: make(chan string, 10),
|
||||||
TelemetryCh: make(chan telemetry.TelemetryData, 1),
|
TelemetryCh: make(chan telemetry.TelemetryData, 1),
|
||||||
Run: false,
|
Run: false,
|
||||||
StreamView: streamView,
|
StreamView: streamView,
|
||||||
isRunning: false,
|
|
||||||
}
|
}
|
||||||
|
|
||||||
ctrl.registerHooks()
|
ctrl.registerHooks()
|
||||||
// ctrl.subscribeListeners()
|
|
||||||
|
|
||||||
return ctrl
|
return ctrl
|
||||||
}
|
}
|
||||||
@@ -96,11 +94,11 @@ func (sc *StreamingCtrl) registerHooks() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (sc *StreamingCtrl) StartStop() {
|
func (sc *StreamingCtrl) StartStop() {
|
||||||
if sc.isRunning {
|
if sc.TelemServ.IsStreaming() {
|
||||||
slog.Info("stopping stream")
|
slog.Info("stopping stream")
|
||||||
|
|
||||||
sc.TelemServ.StopStream()
|
sc.TelemServ.StopStream()
|
||||||
sc.Service.StopStream()
|
sc.DevService.StopStream()
|
||||||
|
|
||||||
sc.isRunning = false
|
sc.isRunning = false
|
||||||
return
|
return
|
||||||
@@ -110,9 +108,9 @@ func (sc *StreamingCtrl) StartStop() {
|
|||||||
// NOTE:
|
// NOTE:
|
||||||
// Subscribe the only existing device - needs to be discovered by now
|
// Subscribe the only existing device - needs to be discovered by now
|
||||||
slog.Debug("setting the data stream for device servie")
|
slog.Debug("setting the data stream for device servie")
|
||||||
sc.Service.SetTelemetryChannel(sc.TelemServ.SubscribeListener("DeviceService", 1))
|
sc.DevService.SetTelemetryChannel(sc.TelemServ.SubscribeListener("DeviceService", 1))
|
||||||
|
|
||||||
dev, err := sc.Service.GetDevice(uidevice.NAME)
|
dev, err := sc.DevService.GetPeripheral(uidevice.NAME)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
if uiDev, ok := dev.(*uidevice.UIDevice); ok {
|
if uiDev, ok := dev.(*uidevice.UIDevice); ok {
|
||||||
sc.TelemetryCh = uiDev.DataChannel()
|
sc.TelemetryCh = uiDev.DataChannel()
|
||||||
@@ -121,7 +119,7 @@ func (sc *StreamingCtrl) StartStop() {
|
|||||||
}
|
}
|
||||||
|
|
||||||
slog.Debug("starting services")
|
slog.Debug("starting services")
|
||||||
sc.Service.StartStream()
|
sc.DevService.StartStream()
|
||||||
sc.TelemServ.StartStream()
|
sc.TelemServ.StartStream()
|
||||||
|
|
||||||
sc.isRunning = true
|
sc.isRunning = true
|
||||||
@@ -161,24 +159,11 @@ func (sc *StreamingCtrl) updateStream() {
|
|||||||
// Performance reasoning: this is not used during the high frequency data transmission
|
// Performance reasoning: this is not used during the high frequency data transmission
|
||||||
// so we can get away with using a map for convenience here
|
// so we can get away with using a map for convenience here
|
||||||
func (sc *StreamingCtrl) SetInternalState() {
|
func (sc *StreamingCtrl) SetInternalState() {
|
||||||
// Acquire the cdashdisplay
|
fields := sc.TelemServ.SubscribeToFields()
|
||||||
// displayIF, err := sc.Service.GetDevice(cdashdisplay.NAME)
|
|
||||||
// if err != nil {
|
|
||||||
// sc.Messages <- "failed to get " + cdashdisplay.NAME
|
|
||||||
// return
|
|
||||||
// }
|
|
||||||
// display, ok := displayIF.(*cdashdisplay.CDashDisplay)
|
|
||||||
// if !ok {
|
|
||||||
// sc.Messages <- "failed to acquire " + cdashdisplay.NAME
|
|
||||||
// return
|
|
||||||
// }
|
|
||||||
// ---
|
|
||||||
|
|
||||||
sc.TelemServ.SubscribeToFields()
|
|
||||||
|
|
||||||
// sc.Messages <- fmt.Sprintf("Subscribed Fields: %+v [%d]\n", fields, len(fields))
|
// sc.Messages <- fmt.Sprintf("Subscribed Fields: %+v [%d]\n", fields, len(fields))
|
||||||
// Should I update this?
|
// Should I update this?
|
||||||
sc.Messages <- fmt.Sprintf("Subscribed to fields\n")
|
sc.Messages <- fmt.Sprintf("Subscribed to fields: %+v", fields)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sc *StreamingCtrl) listenToUIStream() {
|
func (sc *StreamingCtrl) listenToUIStream() {
|
||||||
@@ -190,8 +175,6 @@ func (sc *StreamingCtrl) listenToUIStream() {
|
|||||||
}
|
}
|
||||||
isDrawing.Store(true)
|
isDrawing.Store(true)
|
||||||
|
|
||||||
// sc.Logger.Debug("got data", "data", msg)
|
|
||||||
|
|
||||||
// Capture locally
|
// Capture locally
|
||||||
telemetryMsg := msg
|
telemetryMsg := msg
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user