From dda300ce6b54f3b89afb78efa353de54722e2c9a Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Thu, 17 Sep 2026 23:55:54 +0100 Subject: [PATCH 1/3] decoupled the winIDs from the telemetry data hell yeah, brother! --- devices/cdashdisplay/display.go | 49 ++++++++-- devices/cdashdisplay/transportPackets.go | 77 +++++++++++++++- main.go | 4 +- providers/beamng/beamng.go | 4 +- providers/iracing/iracing.go | 4 +- services/telemetry.go | 2 + telemetry/data.go | 109 ++++++++--------------- telemetry/telemetry.go | 6 +- 8 files changed, 164 insertions(+), 91 deletions(-) diff --git a/devices/cdashdisplay/display.go b/devices/cdashdisplay/display.go index 0c98e8d..dcb6489 100644 --- a/devices/cdashdisplay/display.go +++ b/devices/cdashdisplay/display.go @@ -8,6 +8,7 @@ import ( "log/slog" "os" "path" + "sync" "time" helper "esdi/helpers" @@ -106,8 +107,10 @@ func NewCDashState() *CDashState { } type CDashDisplay struct { - WT *communication.WalkieTalkie - State *CDashState + WT *communication.WalkieTalkie + State *CDashState + fieldToWindows map[telemetry.FieldID][]int16 + bufPool sync.Pool } // Connect will try to find and connect to the CDashDisplay @@ -120,14 +123,42 @@ func NewCDashDisplay() (*CDashDisplay, error) { } return &CDashDisplay{ - WT: p, - State: NewCDashState(), + 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 { diff --git a/devices/cdashdisplay/transportPackets.go b/devices/cdashdisplay/transportPackets.go index 811684b..ad85aff 100644 --- a/devices/cdashdisplay/transportPackets.go +++ b/devices/cdashdisplay/transportPackets.go @@ -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 +} diff --git a/main.go b/main.go index 7507664..f798f9f 100644 --- a/main.go +++ b/main.go @@ -10,7 +10,6 @@ import ( "esdi/cmd" "esdi/config" - "esdi/telemetry" "github.com/arl/statsviz" ) @@ -44,7 +43,7 @@ func initApplication() { } // Setting up some internal data structures - telemetry.Init() + // telemetry.Init() } func setupLogger() error { @@ -56,6 +55,7 @@ func setupLogger() error { logger := slog.New( slog.NewTextHandler(output, &slog.HandlerOptions{ Level: slog.LevelDebug, + // AddSource: true, // NOTE: this may have some performance impacts, disable it for prod }), ) diff --git a/providers/beamng/beamng.go b/providers/beamng/beamng.go index 70e90ed..2f6326f 100644 --- a/providers/beamng/beamng.go +++ b/providers/beamng/beamng.go @@ -125,8 +125,8 @@ 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 { + // b.data.Values[id].IDs = append(b.data.Values[id].IDs, winID) switch id { case telemetry.RPMStateColour: diff --git a/providers/iracing/iracing.go b/providers/iracing/iracing.go index 2c57f6e..9bae148 100644 --- a/providers/iracing/iracing.go +++ b/providers/iracing/iracing.go @@ -229,8 +229,8 @@ 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 { + // i.data.Values[id].IDs = append(i.data.Values[id].IDs, winID) switch id { case telemetry.RPMStateColour: diff --git a/services/telemetry.go b/services/telemetry.go index 9dd7c53..253ec73 100644 --- a/services/telemetry.go +++ b/services/telemetry.go @@ -206,8 +206,10 @@ 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 + // REVAMP for _, dev := range t.devService.Devices { fields := dev.RequiredFields() + t.logger.Debug("requested fields", "fields", fields) t.activeProvider.Subscribe(fields) } } diff --git a/telemetry/data.go b/telemetry/data.go index c4f5ae0..45051f4 100644 --- a/telemetry/data.go +++ b/telemetry/data.go @@ -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,12 +222,12 @@ 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 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] @@ -286,25 +249,25 @@ 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 -} +// 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 +// } diff --git a/telemetry/telemetry.go b/telemetry/telemetry.go index 016e3e1..673b5c3 100644 --- a/telemetry/telemetry.go +++ b/telemetry/telemetry.go @@ -1,5 +1,5 @@ package telemetry -func Init() { - initFieldNamesMap() -} +// func Init() { +// initFieldNamesMap() +// } From c888dc54361fee8678068813e3fba4bf6d4d5457 Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Fri, 18 Sep 2026 00:05:50 +0100 Subject: [PATCH 2/3] remove commented code and made the subscription be based on a list --- devices/cdashdisplay/info.go | 9 +++++---- devices/uidevice/device.go | 11 ++--------- main.go | 4 ---- peripheral/peripheral.go | 2 +- providers/beamng/beamng.go | 4 +--- providers/iracing/iracing.go | 4 +--- telemetry/data.go | 30 ------------------------------ telemetry/provider.go | 2 +- telemetry/telemetry.go | 4 ---- 9 files changed, 11 insertions(+), 59 deletions(-) diff --git a/devices/cdashdisplay/info.go b/devices/cdashdisplay/info.go index 14cbbe4..1dfb8d2 100644 --- a/devices/cdashdisplay/info.go +++ b/devices/cdashdisplay/info.go @@ -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 diff --git a/devices/uidevice/device.go b/devices/uidevice/device.go index d228982..c08fd85 100644 --- a/devices/uidevice/device.go +++ b/devices/uidevice/device.go @@ -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 } diff --git a/main.go b/main.go index f798f9f..7a54f13 100644 --- a/main.go +++ b/main.go @@ -41,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 { @@ -55,7 +52,6 @@ func setupLogger() error { logger := slog.New( slog.NewTextHandler(output, &slog.HandlerOptions{ Level: slog.LevelDebug, - // AddSource: true, // NOTE: this may have some performance impacts, disable it for prod }), ) diff --git a/peripheral/peripheral.go b/peripheral/peripheral.go index bdb035d..f7b31e5 100644 --- a/peripheral/peripheral.go +++ b/peripheral/peripheral.go @@ -18,7 +18,7 @@ const ( type Peripheral interface { Name() string SendData(*telemetry.TelemetryData) - RequiredFields() map[int16]telemetry.FieldID + RequiredFields() []telemetry.FieldID } type PeripheralDeviceClerk struct { diff --git a/providers/beamng/beamng.go b/providers/beamng/beamng.go index 2f6326f..816635d 100644 --- a/providers/beamng/beamng.go +++ b/providers/beamng/beamng.go @@ -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))) @@ -126,8 +126,6 @@ func (b *BeamNG) Subscribe(requestFields map[int16]telemetry.FieldID) { pendingBinds := make([]telemetry.FieldID, telemetry.MaxFields) for _, id := range requestFields { - // b.data.Values[id].IDs = append(b.data.Values[id].IDs, winID) - switch id { case telemetry.RPMStateColour: b.data.VirtualBinds = append(b.data.VirtualBinds, telemetry.NewRPMLights()) diff --git a/providers/iracing/iracing.go b/providers/iracing/iracing.go index 9bae148..9620c3b 100644 --- a/providers/iracing/iracing.go +++ b/providers/iracing/iracing.go @@ -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)) @@ -230,8 +230,6 @@ func (i *IRacing) Subscribe(requestFields map[int16]telemetry.FieldID) { pendingBinds := make([]telemetry.FieldID, 0, telemetry.MaxFields) for _, id := range requestFields { - // i.data.Values[id].IDs = append(i.data.Values[id].IDs, winID) - switch id { case telemetry.RPMStateColour: i.data.VirtualBinds = append(i.data.VirtualBinds, telemetry.NewRPMLights()) diff --git a/telemetry/data.go b/telemetry/data.go index 45051f4..166b17c 100644 --- a/telemetry/data.go +++ b/telemetry/data.go @@ -222,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 @@ -248,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 -// } diff --git a/telemetry/provider.go b/telemetry/provider.go index 93a405a..43c8370 100644 --- a/telemetry/provider.go +++ b/telemetry/provider.go @@ -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() diff --git a/telemetry/telemetry.go b/telemetry/telemetry.go index 673b5c3..073f40d 100644 --- a/telemetry/telemetry.go +++ b/telemetry/telemetry.go @@ -1,5 +1 @@ package telemetry - -// func Init() { -// initFieldNamesMap() -// } From bd10d6b4d644103ca77231db78235eca17fe0144 Mon Sep 17 00:00:00 2001 From: Eduardo Silva Date: Fri, 18 Sep 2026 00:07:27 +0100 Subject: [PATCH 3/3] subscription fix. should subscribe too fields at once --- services/telemetry.go | 20 +++++++++++++------- 1 file changed, 13 insertions(+), 7 deletions(-) diff --git a/services/telemetry.go b/services/telemetry.go index 253ec73..d4bb5db 100644 --- a/services/telemetry.go +++ b/services/telemetry.go @@ -7,6 +7,7 @@ import ( "time" "esdi/providers" + "esdi/telemetry" telem "esdi/telemetry" ) @@ -203,15 +204,20 @@ 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 - // REVAMP + seen := make(map[telemetry.FieldID]struct{}) + var allFields []telemetry.FieldID + for _, dev := range t.devService.Devices { - fields := dev.RequiredFields() - t.logger.Debug("requested fields", "fields", fields) - 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() {