diff --git a/docs/outputs.md b/docs/outputs.md index 6b64757..e749ce2 100644 --- a/docs/outputs.md +++ b/docs/outputs.md @@ -1,7 +1,9 @@ # Outputs Mapd outputs data using two different cereal messages. MapdOut is the primary -output of mapd and is output at a rate of 20 hz. MapdExtendedOut contains data +output of mapd and is output at a rate of 20 hz. It repeats the latest fully +calculated output while new map calculations are in progress, so its publish +cadence does not indicate the calculation rate. MapdExtendedOut contains data that is needed in less realtime and thus is output at a rate of 1hz. ## MapdOut diff --git a/main.go b/main.go index 41aaccf..cfdc765 100644 --- a/main.go +++ b/main.go @@ -2,10 +2,14 @@ package main import ( "log/slog" + "sync/atomic" "time" + "capnproto.org/go/capnp/v3" "github.com/pkg/errors" "pfeifer.dev/mapd/cereal" + "pfeifer.dev/mapd/cereal/custom" + "pfeifer.dev/mapd/cereal/log" "pfeifer.dev/mapd/cli" "pfeifer.dev/mapd/maps" m "pfeifer.dev/mapd/math" @@ -14,6 +18,30 @@ import ( const mapLoadRetryDelay = time.Second +func publishMapdOut(publisher *cereal.Publisher[custom.MapdOut], latest *atomic.Pointer[capnp.Message], stop <-chan struct{}, stopped chan<- struct{}) { + ticker := time.NewTicker(ms.LOOP_DELAY) + defer ticker.Stop() + defer close(stopped) + + for { + message := latest.Load() + event, err := log.ReadRootEvent(message) + if err == nil { + event.SetLogMonoTime(cereal.GetTime()) + err = publisher.Send(message) + } + if err != nil { + slog.Error("Failed to send update", "error", err) + } + + select { + case <-ticker.C: + case <-stop: + return + } + } +} + func main() { ms.Settings.Default() // set defaults so settings not already in param are defaulted settingsLoaded := ms.Settings.Load() // try loading settings before cli @@ -36,7 +64,6 @@ func main() { pub := cereal.NewPublisher("mapdOut", cereal.MapdOutCreator) defer pub.Pub.Msgq.Close() - state.Publisher = &pub sub := cereal.NewSubscriber("mapdIn", cereal.MapdInReader, false, false) defer sub.Sub.Msgq.Close() @@ -56,18 +83,21 @@ func main() { selfdriveState := cereal.NewSubscriber("selfdriveState", cereal.SelfdriveStateReader, true, ms.Settings.SubscriberSettings.ShadowSelfdriveState) defer selfdriveState.Sub.Msgq.Close() - lastLoopTime := time.Now() + latestMapdOut := atomic.Pointer[capnp.Message]{} + latestMapdOut.Store(state.buildMapdOut(&pub)) + stopMapdOut := make(chan struct{}) + mapdOutStopped := make(chan struct{}) + go publishMapdOut(&pub, &latestMapdOut, stopMapdOut, mapdOutStopped) + defer func() { + close(stopMapdOut) + <-mapdOutStopped + }() for { - lastLoopDuration := time.Since(lastLoopTime) - time.Sleep(max(ms.LOOP_DELAY - lastLoopDuration, 0)) - lastLoopTime = time.Now() + // Build before handoff because State contains unsynchronized lazy getters. + latestMapdOut.Store(state.buildMapdOut(&pub)) - err := state.Send() // send beginning of each loop to ensure it happens at the correct rate - if err != nil { - slog.Error("Failed to send update", "error", err) - } - err = extendedState.Send() // this send is internally rate limited to 1 hz + err := extendedState.Send() // this send is internally rate limited to 1 hz if err != nil { slog.Error("Failed to send extended update", "error", err) } @@ -136,6 +166,5 @@ func main() { state.TargetVelocities = GetTargetVelocities(state.Curvatures, state.TargetVelocities) } - // send at beginning of next loop } } diff --git a/state.go b/state.go index 1e7915f..eb8672e 100644 --- a/state.go +++ b/state.go @@ -1,6 +1,7 @@ package main import ( + "capnproto.org/go/capnp/v3" "pfeifer.dev/mapd/cereal" "pfeifer.dev/mapd/cereal/car" "pfeifer.dev/mapd/cereal/custom" @@ -10,7 +11,6 @@ import ( ) type State struct { - Publisher *cereal.Publisher[custom.MapdOut] Data maps.Offline Car CarState CurrentWay CurrentWay @@ -65,8 +65,8 @@ func (s *State) UpdateCarState(carData car.CarState) { s.SpeedLimit.Update(s.CurrentWay, s.Car) } -func (s *State) Send() error { - msg, output := s.Publisher.NewMessage(true) +func (s *State) buildMapdOut(publisher *cereal.Publisher[custom.MapdOut]) *capnp.Message { + msg, output := publisher.NewMessage(true) id := s.CurrentWay.Way.Id() output.SetWayId(id) @@ -120,5 +120,5 @@ func (s *State) Send() error { output.SetWaySelectionType(s.CurrentWay.SelectionType) output.SetSpeedLimitAccepted(ms.Settings.SpeedLimitAccepted()) - return s.Publisher.Send(msg) + return msg }