decoupled the winIDs from the telemetry data #11
@@ -8,6 +8,7 @@ import (
|
||||
"log/slog"
|
||||
"os"
|
||||
"path"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
helper "esdi/helpers"
|
||||
@@ -108,6 +109,8 @@ func NewCDashState() *CDashState {
|
||||
type CDashDisplay struct {
|
||||
WT *communication.WalkieTalkie
|
||||
State *CDashState
|
||||
fieldToWindows map[telemetry.FieldID][]int16
|
||||
bufPool sync.Pool
|
||||
}
|
||||
|
||||
// Connect will try to find and connect to the CDashDisplay
|
||||
@@ -122,12 +125,40 @@ func NewCDashDisplay() (*CDashDisplay, error) {
|
||||
return &CDashDisplay{
|
||||
WT: p,
|
||||
State: NewCDashState(),
|
||||
fieldToWindows: make(map[telemetry.FieldID][]int16),
|
||||
bufPool: sync.Pool{
|
||||
New: func() any {
|
||||
b := make([]byte, 0, telemetry.MaxFields*8)
|
||||
return &b
|
||||
},
|
||||
},
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (d *CDashDisplay) SendCommand() {
|
||||
}
|
||||
|
||||
func (d *CDashDisplay) RegisterFieldMapping(fieldID telemetry.FieldID, winID int16) {
|
||||
d.fieldToWindows[fieldID] = append(d.fieldToWindows[fieldID], winID)
|
||||
}
|
||||
|
||||
func (d *CDashDisplay) UnregisterFieldMapping(winID int16) {
|
||||
for fieldID, windows := range d.fieldToWindows {
|
||||
updated := windows[:0]
|
||||
for _, w := range windows {
|
||||
if w != winID {
|
||||
updated = append(updated, w)
|
||||
}
|
||||
|
||||
if len(updated) == 0 {
|
||||
delete(d.fieldToWindows, fieldID)
|
||||
} else {
|
||||
d.fieldToWindows[fieldID] = updated
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (d *CDashDisplay) CreateWindow(win *DesktopUIWindow) (*DesktopUIWindow, error) {
|
||||
bytes, err := helper.StructToBytes(win.UIWindow)
|
||||
if err != nil {
|
||||
@@ -146,6 +177,9 @@ func (d *CDashDisplay) CreateWindow(win *DesktopUIWindow) (*DesktopUIWindow, err
|
||||
slog.Info(fmt.Sprintf("Recived ID message: %v", wID))
|
||||
|
||||
d.State.Layout.AddWindow(win)
|
||||
if fieldID, ok := telemetry.GetFieldID(win.UIData.TelemetryField); ok {
|
||||
d.RegisterFieldMapping(fieldID, win.UIData.IDX)
|
||||
}
|
||||
|
||||
return win, nil
|
||||
}
|
||||
@@ -177,6 +211,8 @@ func (d *CDashDisplay) UpdateWindow(win *DesktopUIWindow) error {
|
||||
// Yeah, same address as suspected
|
||||
// I can't think about it right now. I'll think about that tomorrow
|
||||
|
||||
// TODO: need to update the field mappings here!
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -199,7 +235,8 @@ func (d *CDashDisplay) DestroyWindow(wID int16) error {
|
||||
return err
|
||||
}
|
||||
|
||||
// NODE: add this
|
||||
// NOTE: add this
|
||||
d.UnregisterFieldMapping(wID)
|
||||
d.State.Layout.RemoveWindow(wID)
|
||||
|
||||
return nil
|
||||
@@ -357,7 +394,7 @@ func (d *CDashDisplay) UnloadLayout() error {
|
||||
}
|
||||
|
||||
func (d *CDashDisplay) SendData(data *telemetry.TelemetryData) {
|
||||
packet := data.Pack()
|
||||
packet := d.encodePacket(data)
|
||||
|
||||
bytes, err := helper.StructToBytes(packet)
|
||||
if err != nil {
|
||||
|
||||
@@ -15,12 +15,13 @@ func (cds *CDashDisplay) Name() string {
|
||||
return NAME
|
||||
}
|
||||
|
||||
func (cds *CDashDisplay) RequiredFields() map[int16]telemetry.FieldID {
|
||||
fields := make(map[int16]telemetry.FieldID, len(cds.State.Layout.Windows))
|
||||
func (cds *CDashDisplay) RequiredFields() []telemetry.FieldID {
|
||||
fields := make([]telemetry.FieldID, 0, len(cds.State.Layout.Windows))
|
||||
|
||||
for _, w := range cds.State.Layout.Windows {
|
||||
fieldID, _ := telemetry.GetFieldID(w.UIData.TelemetryField)
|
||||
fields[w.UIData.IDX] = fieldID
|
||||
if fieldID, ok := telemetry.GetFieldID(w.UIData.TelemetryField); ok {
|
||||
fields = append(fields, fieldID)
|
||||
}
|
||||
}
|
||||
|
||||
return fields
|
||||
|
||||
@@ -1,5 +1,11 @@
|
||||
package cdashdisplay
|
||||
|
||||
import (
|
||||
"math"
|
||||
|
||||
"esdi/telemetry"
|
||||
)
|
||||
|
||||
// In this file we will place all structs that are 1:1 representation of the
|
||||
// types in the device transport layer ->
|
||||
|
||||
@@ -65,6 +71,71 @@ type UIWindowUpdatePacket struct {
|
||||
Window UIWindow
|
||||
}
|
||||
|
||||
// func getUIWindowDTO(w *DesktopWindowData) *UIWindow {
|
||||
//
|
||||
// }
|
||||
// TODO: this is here because this is the only device I use like this
|
||||
// once I add another serial device I will move this somewhere else that can
|
||||
// be reused by all devices but still is decoupled from the telemetry package
|
||||
func (cds *CDashDisplay) encodePacket(td *telemetry.TelemetryData) []byte {
|
||||
bufPtr := cds.bufPool.Get().(*[]byte)
|
||||
buf := (*bufPtr)[:0]
|
||||
|
||||
for fieldID, windowIDs := range cds.fieldToWindows {
|
||||
if int(fieldID) >= len(td.Values) || len(windowIDs) == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
tf := &td.Values[fieldID]
|
||||
|
||||
for _, winID := range windowIDs {
|
||||
buf = packField(winID, tf, buf)
|
||||
}
|
||||
}
|
||||
|
||||
// We have to copy here because we have to return the buffer
|
||||
result := make([]byte, len(buf))
|
||||
copy(result, buf)
|
||||
|
||||
cds.bufPool.Put(&buf)
|
||||
|
||||
return result
|
||||
}
|
||||
|
||||
// Pack will pack this current TelemetryField into bytes to send over the wire
|
||||
// Format:
|
||||
// 0x00 - Field ID
|
||||
// 0x00 |
|
||||
// 0x01 - DataType
|
||||
// 0x02 - if its a (u)int8
|
||||
// or
|
||||
// 0x02 - if its a (u)int16 - first byte
|
||||
// 0x02 - if its a (u)int16 - second byte
|
||||
// or
|
||||
// 0x02 - str len max is 255 chars
|
||||
// [0x02] - str
|
||||
func packField(winID int16, tf *telemetry.TelemetryField, dest []byte) []byte {
|
||||
// NOTE: maybe we can have a pool of these so we don't have to create them here
|
||||
// or whatever
|
||||
dest = append(dest, uint8(winID), uint8(winID>>8))
|
||||
dest = append(dest, uint8(tf.Type))
|
||||
|
||||
switch tf.Type {
|
||||
case telemetry.DataTypeINT8, telemetry.DataTypeUINT8, telemetry.DataTypeCHAR:
|
||||
dest = append(dest, uint8(tf.Raw))
|
||||
case telemetry.DataTypeINT16, telemetry.DataTypeUINT16:
|
||||
dest = append(dest, uint8(tf.Raw), uint8(tf.Raw>>8))
|
||||
case telemetry.DataTypeINT32, telemetry.DataTypeUINT32:
|
||||
dest = append(dest, uint8(tf.Raw), uint8(tf.Raw>>8), uint8(tf.Raw>>16), uint8(tf.Raw>>24))
|
||||
case telemetry.DataTypeINT64, telemetry.DataTypeUINT64:
|
||||
dest = append(
|
||||
dest, uint8(tf.Raw), uint8(tf.Raw>>8), uint8(tf.Raw>>16),
|
||||
uint8(tf.Raw>>24), uint8(tf.Raw>>32), uint8(tf.Raw>>40), uint8(tf.Raw>>48),
|
||||
uint8(tf.Raw>>56),
|
||||
)
|
||||
case telemetry.DataTypeSTRING:
|
||||
l := min(len(tf.Str), math.MaxUint8)
|
||||
|
||||
dest = append(dest, uint8(l))
|
||||
dest = append(dest, tf.Str[:l]...)
|
||||
}
|
||||
|
||||
return dest
|
||||
}
|
||||
|
||||
@@ -37,17 +37,10 @@ func (uid *UIDevice) DataChannel() <-chan telemetry.TelemetryData {
|
||||
return uid.dataChan
|
||||
}
|
||||
|
||||
func (uid *UIDevice) RequiredFields() map[int16]telemetry.FieldID {
|
||||
subscribeTo := []telemetry.FieldID{
|
||||
func (uid *UIDevice) RequiredFields() []telemetry.FieldID {
|
||||
return []telemetry.FieldID{
|
||||
telemetry.Speed,
|
||||
telemetry.Gear,
|
||||
telemetry.RPM,
|
||||
}
|
||||
|
||||
fields := make(map[int16]telemetry.FieldID, len(subscribeTo))
|
||||
for k, field := range subscribeTo {
|
||||
fields[int16(k)] = field
|
||||
}
|
||||
|
||||
return fields
|
||||
}
|
||||
|
||||
@@ -10,7 +10,6 @@ import (
|
||||
|
||||
"esdi/cmd"
|
||||
"esdi/config"
|
||||
"esdi/telemetry"
|
||||
|
||||
"github.com/arl/statsviz"
|
||||
)
|
||||
@@ -42,9 +41,6 @@ func initApplication() {
|
||||
slog.Error(fmt.Sprintf("failed to setup metrics server: %+v", err))
|
||||
}
|
||||
}
|
||||
|
||||
// Setting up some internal data structures
|
||||
telemetry.Init()
|
||||
}
|
||||
|
||||
func setupLogger() error {
|
||||
|
||||
@@ -18,7 +18,7 @@ const (
|
||||
type Peripheral interface {
|
||||
Name() string
|
||||
SendData(*telemetry.TelemetryData)
|
||||
RequiredFields() map[int16]telemetry.FieldID
|
||||
RequiredFields() []telemetry.FieldID
|
||||
}
|
||||
|
||||
type PeripheralDeviceClerk struct {
|
||||
|
||||
@@ -115,7 +115,7 @@ func (b *BeamNG) Stream() (<-chan telemetry.TelemetryData, error) {
|
||||
return b.streamCh, nil
|
||||
}
|
||||
|
||||
func (b *BeamNG) Subscribe(requestFields map[int16]telemetry.FieldID) {
|
||||
func (b *BeamNG) Subscribe(requestFields []telemetry.FieldID) {
|
||||
// NOTE: document how the Subscribe funtion works
|
||||
slog.Debug(fmt.Sprintf("Len Req: %d\n", len(requestFields)))
|
||||
|
||||
@@ -125,9 +125,7 @@ func (b *BeamNG) Subscribe(requestFields map[int16]telemetry.FieldID) {
|
||||
// we will add their dependencies and the primitives to a slice
|
||||
pendingBinds := make([]telemetry.FieldID, telemetry.MaxFields)
|
||||
|
||||
for winID, id := range requestFields {
|
||||
b.data.Values[id].IDs = append(b.data.Values[id].IDs, winID)
|
||||
|
||||
for _, id := range requestFields {
|
||||
switch id {
|
||||
case telemetry.RPMStateColour:
|
||||
b.data.VirtualBinds = append(b.data.VirtualBinds, telemetry.NewRPMLights())
|
||||
|
||||
@@ -220,7 +220,7 @@ func (i *IRacing) StopStream() {
|
||||
i.streamCancel = nil
|
||||
}
|
||||
|
||||
func (i *IRacing) Subscribe(requestFields map[int16]telemetry.FieldID) {
|
||||
func (i *IRacing) Subscribe(requestFields []telemetry.FieldID) {
|
||||
i.logger.Debug(fmt.Sprintf("Len Req: %d\n", len(requestFields)))
|
||||
|
||||
i.data.ActiveBinds = make([]telemetry.BoundField, 0, len(requestFields))
|
||||
@@ -229,9 +229,7 @@ func (i *IRacing) Subscribe(requestFields map[int16]telemetry.FieldID) {
|
||||
// we will add their dependencies and the primitives to a slice
|
||||
pendingBinds := make([]telemetry.FieldID, 0, telemetry.MaxFields)
|
||||
|
||||
for winID, id := range requestFields {
|
||||
i.data.Values[id].IDs = append(i.data.Values[id].IDs, winID)
|
||||
|
||||
for _, id := range requestFields {
|
||||
switch id {
|
||||
case telemetry.RPMStateColour:
|
||||
i.data.VirtualBinds = append(i.data.VirtualBinds, telemetry.NewRPMLights())
|
||||
|
||||
+13
-5
@@ -7,6 +7,7 @@ import (
|
||||
"time"
|
||||
|
||||
"esdi/providers"
|
||||
"esdi/telemetry"
|
||||
telem "esdi/telemetry"
|
||||
)
|
||||
|
||||
@@ -203,14 +204,21 @@ func (t *TelemetryService) UnsubscribeListener(id string) {
|
||||
}
|
||||
|
||||
func (t *TelemetryService) SubscribeToFields() {
|
||||
// _ = t.devService.SubscribeFields()
|
||||
// TODO: we need to find a way of requesting devices to send all subscribed fields
|
||||
// instead of going through the devices on DevService here
|
||||
seen := make(map[telemetry.FieldID]struct{})
|
||||
var allFields []telemetry.FieldID
|
||||
|
||||
for _, dev := range t.devService.Devices {
|
||||
fields := dev.RequiredFields()
|
||||
t.activeProvider.Subscribe(fields)
|
||||
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")
|
||||
|
||||
+8
-75
@@ -1,12 +1,18 @@
|
||||
package telemetry
|
||||
|
||||
import (
|
||||
"math"
|
||||
"strconv"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
func init() {
|
||||
fieldNameToID = make(map[string]FieldID, MaxFields)
|
||||
for id, name := range FieldNames {
|
||||
fieldNameToID[name] = FieldID(id)
|
||||
}
|
||||
}
|
||||
|
||||
// NOTE: allow the user to create custom data things. For example, iRacing provides
|
||||
// multiple surface temps, but I guess the user doesn't want all of them at once.
|
||||
// allow him to make something that allows some data transformation to occur.
|
||||
@@ -64,7 +70,7 @@ const (
|
||||
// NOTE: we can optimize this via a special command that says a given piece of data
|
||||
// is for multiple targets
|
||||
type TelemetryField struct {
|
||||
IDs []int16 // Identification for the serial device
|
||||
// IDs []int16 // Identification for the serial device
|
||||
Type DataType
|
||||
Raw uint64
|
||||
Str string // Only to be used with DataTypeSTRING
|
||||
@@ -75,49 +81,6 @@ func (tf *TelemetryField) Unused() {
|
||||
tf.Raw = uint64('-')
|
||||
}
|
||||
|
||||
// Pack will pack this current TelemetryField into bytes to send over the wire
|
||||
// Format:
|
||||
// 0x00 - Field ID
|
||||
// 0x00 |
|
||||
// 0x01 - DataType
|
||||
// 0x02 - if its a (u)int8
|
||||
// or
|
||||
// 0x02 - if its a (u)int16 - first byte
|
||||
// 0x02 - if its a (u)int16 - second byte
|
||||
// or
|
||||
// 0x02 - str len max is 255 chars
|
||||
// [0x02] - str
|
||||
func (tf *TelemetryField) Pack(dest []byte) []byte {
|
||||
// NOTE: maybe we can have a pool of these so we don't have to create them here
|
||||
// or whatever
|
||||
for _, id := range tf.IDs {
|
||||
dest = append(dest, uint8(id), uint8(id>>8))
|
||||
dest = append(dest, uint8(tf.Type))
|
||||
|
||||
switch tf.Type {
|
||||
case DataTypeINT8, DataTypeUINT8, DataTypeCHAR:
|
||||
dest = append(dest, uint8(tf.Raw))
|
||||
case DataTypeINT16, DataTypeUINT16:
|
||||
dest = append(dest, uint8(tf.Raw), uint8(tf.Raw>>8))
|
||||
case DataTypeINT32, DataTypeUINT32:
|
||||
dest = append(dest, uint8(tf.Raw), uint8(tf.Raw>>8), uint8(tf.Raw>>16), uint8(tf.Raw>>24))
|
||||
case DataTypeINT64, DataTypeUINT64:
|
||||
dest = append(
|
||||
dest, uint8(tf.Raw), uint8(tf.Raw>>8), uint8(tf.Raw>>16),
|
||||
uint8(tf.Raw>>24), uint8(tf.Raw>>32), uint8(tf.Raw>>40), uint8(tf.Raw>>48),
|
||||
uint8(tf.Raw>>56),
|
||||
)
|
||||
case DataTypeSTRING:
|
||||
l := min(len(tf.Str), math.MaxUint8)
|
||||
|
||||
dest = append(dest, uint8(l))
|
||||
dest = append(dest, tf.Str[:l]...)
|
||||
}
|
||||
}
|
||||
|
||||
return dest
|
||||
}
|
||||
|
||||
func (tf *TelemetryField) String() string {
|
||||
switch tf.Type {
|
||||
case DataTypeSTRING:
|
||||
@@ -259,13 +222,6 @@ func GetFieldName(id FieldID) string {
|
||||
|
||||
var fieldNameToID map[string]FieldID
|
||||
|
||||
func initFieldNamesMap() {
|
||||
fieldNameToID = make(map[string]FieldID, MaxFields)
|
||||
for id, name := range FieldNames {
|
||||
fieldNameToID[name] = FieldID(id)
|
||||
}
|
||||
}
|
||||
|
||||
func GetFieldID(name string) (FieldID, bool) {
|
||||
id, ok := fieldNameToID[name]
|
||||
return id, ok
|
||||
@@ -285,26 +241,3 @@ type TelemetryData struct {
|
||||
func NewTelemetryData() *TelemetryData {
|
||||
return &TelemetryData{}
|
||||
}
|
||||
|
||||
func (td *TelemetryData) Pack() []byte {
|
||||
bufPtr := bufferPool.Get().(*[]byte)
|
||||
buf := (*bufPtr)[:0]
|
||||
|
||||
// for _, bind := range td.ActiveBinds {
|
||||
// buf = td.Values[bind.ID].Pack(buf)
|
||||
// }
|
||||
|
||||
for k := range td.Values {
|
||||
if len(td.Values[k].IDs) > 0 {
|
||||
buf = td.Values[k].Pack(buf)
|
||||
}
|
||||
}
|
||||
|
||||
// We have to copy here because we have to return the buffer
|
||||
result := make([]byte, len(buf))
|
||||
copy(result, buf)
|
||||
|
||||
bufferPool.Put(&buf)
|
||||
|
||||
return result
|
||||
}
|
||||
|
||||
@@ -6,7 +6,7 @@ import "time"
|
||||
type TelemetryProvider interface {
|
||||
StopStream()
|
||||
Stream() (<-chan TelemetryData, error)
|
||||
Subscribe(map[int16]FieldID)
|
||||
Subscribe([]FieldID)
|
||||
IsAlive(time.Duration) bool
|
||||
Name() string
|
||||
Close()
|
||||
|
||||
@@ -1,5 +1 @@
|
||||
package telemetry
|
||||
|
||||
func Init() {
|
||||
initFieldNamesMap()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user