broadcaster setup

This commit is contained in:
2026-06-24 12:00:02 +01:00
parent 9c7c6d528f
commit 5b77b5b941
8 changed files with 274 additions and 80 deletions
+13
View File
@@ -42,6 +42,19 @@ other drivers, track conditions and so on so forth
- [ ] Telemetry analysis tool
- [ ] Better user interface
## Debugging
### Freezes:
Using delve:
- Launch Terminal1 with `dlv debug . --headless --listen=:2345 -- tui`
- Launch Terminal2 and connect to with with `dlv connect :2345`
- Type `continue` onto Terminal2 and go to Terminal1 to use the application
normally until it hangs.
- Go back to Terminal2 to and do a `Ctrl+c` to capture the state and then check
what went wrong by looking at the `goroutines` for example.
### Shameless begging
Hey, doesn't hurt to try, its free either way:
[Buy me a coffee!](buymeacoffee.com/ESilva_15)
+4 -1
View File
@@ -65,7 +65,10 @@ func runStatviz() error {
}
go func() {
http.ListenAndServe("localhost:8001", mux)
err := http.ListenAndServe("localhost:8001", mux)
if err != nil {
slog.Error("Failed to set up server for statsviz", "err", err)
}
}()
return nil
+12
View File
@@ -94,6 +94,13 @@ func (i *IRacing) stream(ctx context.Context) {
go func() {
for {
// Explicitly intercpt cancellation
select {
case <-ctx.Done():
return
default:
}
// We start by checking if we do or do not have data available
if !i.isDataAvailable() {
continue
@@ -161,7 +168,12 @@ func (i *IRacing) Stream() (<-chan telem.TelemetryData, error) {
}
func (i *IRacing) StopStream() {
if i.streamCancel == nil {
return
}
i.streamCancel()
i.streamCancel = nil
}
func (i *IRacing) Subscribe(requestFields map[int16]telem.FieldID) {
+13 -1
View File
@@ -2,14 +2,18 @@
package providers
import (
"log/slog"
"esdi/providers/beamng"
"esdi/providers/iracing"
"esdi/telemetry"
)
// Make this be some kind of struct where we can access a function that returns
// the selected provider by its name
type Provider struct {
Name string
Name string
Provider telemetry.TelemetryProvider
}
var Providers = map[string]Provider{
@@ -20,3 +24,11 @@ var Providers = map[string]Provider{
Name: iracing.NAME,
},
}
func NewIRacingProvider(logger *slog.Logger, source string,
telemOut string, yamlOut string,
) telemetry.TelemetryProvider {
provider, _ := iracing.NewIRacingProvider(logger, source, "", "")
return provider
}
+41 -8
View File
@@ -1,6 +1,7 @@
package services
import (
"context"
"fmt"
"log/slog"
"sync/atomic"
@@ -12,11 +13,13 @@ import (
)
type CDashService struct {
Logger *slog.Logger
CDash *cdashdisplay.CDashDisplay
// iRacingTelemetry *IRacingService
Logger *slog.Logger
CDash *cdashdisplay.CDashDisplay
DevClerk *peripheral.PeripheralDeviceClerk
Messages chan string
// Telemetry Channel
streamCancel context.CancelFunc
TelemCh <-chan telemetry.TelemetryData
}
func NewCDashService(logger *slog.Logger) *CDashService {
@@ -97,19 +100,49 @@ func (cds *CDashService) MoveWindow(idx int16, vec *helper.Vector) error {
return nil
}
func (cds *CDashService) StreamData(stream <-chan telemetry.TelemetryData) {
func (cds *CDashService) SetTelemetryChannel(ch <-chan telemetry.TelemetryData) {
cds.TelemCh = ch
}
func (cds *CDashService) StartStream() {
// NOTE: i'm using this pattern a whole lot. Maybe I can create a struct to handle this
var ctx context.Context
ctx, cds.streamCancel = context.WithCancel(context.Background())
go cds.transmit(ctx)
}
func (cds *CDashService) StopStream() {
if cds.streamCancel == nil {
return
}
cds.streamCancel()
cds.streamCancel = nil
}
// INTERNAL
func (cds *CDashService) transmit(ctx context.Context) {
var isSending atomic.Bool
go func() {
for msg := range stream {
for {
select {
case <-ctx.Done():
return
case data, ok := <-cds.TelemCh:
if !ok {
return
}
if isSending.Load() {
continue
}
isSending.Store(true)
cds.CDash.SendData(&msg)
cds.CDash.SendData(&data)
isSending.Store(false)
}
}()
}
}
+132 -35
View File
@@ -1,61 +1,158 @@
package services
import (
"context"
"log/slog"
"sync"
providerir "esdi/providers/iracing"
telemetry "esdi/telemetry"
"esdi/providers"
telem "esdi/telemetry"
)
// TelemetryService will be our base struct to handle telemetry data
// It should hook to a data sink and handle it like iRacing, BeamNG, AC and so on
type TelemetryService struct {
logger *slog.Logger
cdash *CDashService
ActiveProvider telemetry.TelemetryProvider
logger *slog.Logger
cdash *CDashService
// Concurrency protection
mut sync.RWMutex
ativeProvider telem.TelemetryProvider
// Channel for the UI
listeners map[string]chan telem.TelemetryData
uiOutCh chan telem.TelemetryData
cancelForward context.CancelFunc
}
func NewTelemetryService(logger *slog.Logger, cdash *CDashService) *TelemetryService {
return &TelemetryService{
logger: logger,
cdash: cdash,
newService := &TelemetryService{
logger: logger,
cdash: cdash,
uiOutCh: make(chan telem.TelemetryData, 100),
listeners: make(map[string]chan telem.TelemetryData),
}
// Need to instantiate a default provider here
source := "/home/esilva/Desktop/projetos/simracing_peripherals/testTelemetry/gt3_mustang_bathurst.ibt"
firstProvider := providers.NewIRacingProvider(slog.Default(), source, "", "")
newService.SwitchProvider(firstProvider)
return newService
}
// func (t *TelemetryService) GetUIStream() <-chan telem.TelemetryData {
// return t.uiOutCh
// }
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.ativeProvider != nil {
t.dropActiveProvider()
}
// Assign the new provider
t.ativeProvider = newProvider
return nil
}
func (t *TelemetryService) multiplexData(ctx context.Context, dataCh <-chan telem.TelemetryData) {
for {
select {
case <-ctx.Done():
return
case data, ok := <-dataCh:
if !ok {
return
}
t.mut.RLock()
for _, ch := range t.listeners {
select {
case ch <- data:
// Sends data to the subscriber
default:
// Subscriber is full, we just skip ahead. Maybe find a how to add metrics here
}
}
t.mut.RUnlock()
select {
case t.uiOutCh <- data:
default:
// Control latency
}
}
}
}
func (t *TelemetryService) setIRacingProvider() {
path := "/home/esilva/Desktop/projetos/simracing_peripherals/testTelemetry/gt3_mustang_bathurst.ibt"
provider, _ := providerir.NewIRacingProvider(t.logger, path, "", "")
t.ActiveProvider = provider
}
func (t *TelemetryService) SetProvider(provider string) *TelemetryService {
if t.ActiveProvider != nil {
// Gotta do something here to clean up before switching
func (t *TelemetryService) dropActiveProvider() {
if t.cancelForward != nil {
t.cancelForward()
}
switch provider {
case "iRacing":
// Set up iRacing
t.setIRacingProvider()
default:
return nil
}
return t
t.ativeProvider.StopStream()
}
func (t *TelemetryService) StartStream() <-chan telemetry.TelemetryData {
stream, _ := t.ActiveProvider.Stream()
func (t *TelemetryService) SubscribeListener(id string, bufferSize int) <-chan telem.TelemetryData {
t.mut.Lock()
defer t.mut.Unlock()
// NOTE: we don't want to just send data to the cdashdisplay. We want to support
// multiple devices at the same time
t.cdash.StreamData(stream)
// NOTE: is this truly necessary?
// return the channel if it already exists
if ch, exists := t.listeners[id]; exists {
return ch
}
return stream
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(fields map[int16]telem.FieldID) {
t.ativeProvider.Subscribe(fields)
}
func (t *TelemetryService) StartStream() {
slog.Debug("Stream started")
// Start the new stream
simInCh, _ := t.ativeProvider.Stream()
// if err != nil {
// // NOTE
// }
// 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.ActiveProvider.StopStream()
t.mut.Lock()
defer t.mut.Unlock()
if t.cancelForward != nil {
t.cancelForward()
t.cancelForward = nil
}
t.ativeProvider.StopStream()
}
+58 -33
View File
@@ -17,13 +17,14 @@ import (
type StreamingCtrl struct {
*Controller
Service *services.CDashService
StreamView *views.StreamToolView
Messages chan string
Internal chan string
Run bool
OnExit func()
TelemServ *services.TelemetryService
Service *services.CDashService
StreamView *views.StreamToolView
Messages chan string
Internal chan string
TelemetryCh <-chan telemetry.TelemetryData
Run bool
OnExit func()
TelemServ *services.TelemetryService
// Stream State
isRunning bool
@@ -42,21 +43,28 @@ func NewStreamingCtrl(
streamView := views.NewStreamToolView(providerList, config.GetCfg().DefaultSim)
ctrl := &StreamingCtrl{
Controller: base,
Service: serCDash,
TelemServ: serTelem,
Messages: make(chan string, 10),
Internal: make(chan string, 10),
Run: false,
StreamView: streamView,
isRunning: false,
Controller: base,
Service: serCDash,
TelemServ: serTelem,
Messages: make(chan string, 10),
Internal: make(chan string, 10),
TelemetryCh: make(chan telemetry.TelemetryData, 100),
Run: false,
StreamView: streamView,
isRunning: false,
}
ctrl.registerHooks()
ctrl.subscribeListeners()
go ctrl.listenToUIStream()
return ctrl
}
func (sc *StreamingCtrl) subscribeListeners() {
sc.TelemetryCh = sc.TelemServ.SubscribeListener("UI", 50)
}
func (sc *StreamingCtrl) registerHooks() {
sc.StreamView.Options.Form.SetInputCapture(func(ev *tcell.EventKey) *tcell.EventKey {
switch ev.Key() {
@@ -88,30 +96,31 @@ func (sc *StreamingCtrl) registerHooks() {
func (sc *StreamingCtrl) StartStop() {
if sc.isRunning {
slog.Info("stopping stream")
sc.TelemServ.StopStream()
sc.Service.StopStream()
sc.isRunning = false
return
}
stream := sc.TelemServ.StartStream()
// stream is not running, we have to start it now
// NOTE:
// Subscribe the only existing device - needs to be discovered by now
slog.Debug("setting the data stream for cdash")
sc.Service.SetTelemetryChannel(sc.TelemServ.SubscribeListener("cdash", 50))
slog.Debug("starting to stream data again")
sc.Service.StartStream()
slog.Debug("starting the stream")
sc.TelemServ.StartStream()
slog.Debug("setting local control variables")
sc.isRunning = true
var isDrawing atomic.Bool
go func() {
for msg := range stream {
if isDrawing.Load() {
continue
}
isDrawing.Store(true)
sc.App.QueueUpdateDraw(func() {
sc.StreamView.Visualizer.Update(&msg)
isDrawing.Store(false)
})
}
}()
slog.Debug("starting stream")
}
func (sc *StreamingCtrl) parseStreamUpdateForm(form *views.StreamOptionsView) (*models.StreamOptions, error) {
@@ -153,7 +162,23 @@ func (sc *StreamingCtrl) SetInternalState() {
fields[w.UIData.IDX] = fieldID
}
sc.TelemServ.ActiveProvider.Subscribe(fields)
sc.TelemServ.SubscribeToFields(fields)
sc.Messages <- fmt.Sprintf("Subscribed Fields: %+v [%d]\n", fields, len(fields))
}
func (sc *StreamingCtrl) listenToUIStream() {
var isDrawing atomic.Bool
for msg := range sc.TelemetryCh {
if isDrawing.Load() {
continue
}
isDrawing.Store(true)
sc.App.QueueUpdateDraw(func() {
sc.StreamView.Visualizer.Update(&msg)
isDrawing.Store(false)
})
}
}
+1 -2
View File
@@ -24,8 +24,7 @@ func NewControlPanel(logger *slog.Logger) *ControlPanel {
}
devService := services.NewCDashService(logger)
telemService := services.NewTelemetryService(logger, devService).
SetProvider("iRacing")
telemService := services.NewTelemetryService(logger, devService)
if telemService == nil {
panic("failed to create the telemetry service")
}