#pragma once #include #include #include #include #include #include #include #include #include #include #include #include #include "cereal/messaging/messaging.h" #include "tools/cabana/core/can_data.h" #include "tools/cabana/core/observable.h" #include "tools/cabana/dbc/dbcmanager.h" #include "tools/cabana/utils/util.h" #include "tools/replay/util.h" class AbstractStream { public: AbstractStream(); virtual ~AbstractStream() = default; virtual void start() = 0; virtual bool liveStreaming() const { return true; } virtual void seekTo(double ts) {} virtual std::string routeName() const = 0; virtual std::string carFingerprint() const { return ""; } virtual std::chrono::system_clock::time_point beginDateTime() const { return {}; } virtual uint64_t beginMonoTime() const { return 0; } virtual double minSeconds() const { return 0; } virtual double maxSeconds() const { return 0; } virtual void setSpeed(float speed) {} virtual double getSpeed() { return 1; } virtual bool isPaused() const { return false; } virtual void pause(bool pause) {} void setTimeRange(const std::optional> &range); const std::optional> &timeRange() const { return time_range_; } inline double currentSec() const { return current_sec_; } inline uint64_t toMonoTime(double sec) const { return beginMonoTime() + std::max(sec, 0.0) * 1e9; } inline double toSeconds(uint64_t mono_time) const { return std::max(0.0, (mono_time - beginMonoTime()) / 1e9); } inline const std::unordered_map &lastMessages() const { return last_msgs; } bool isMessageActive(const MessageId &id) const; inline const MessageEventsMap &eventsMap() const { return events_; } inline const std::vector &allEvents() const { return all_events_; } const CanData &lastMessage(const MessageId &id) const; const std::vector &events(const MessageId &id) const; std::pair eventsInRange(const MessageId &id, std::optional> time_range) const; size_t suppressHighlighted(); void clearSuppressed(); void suppressDefinedSignals(bool suppress); Observable<> paused; Observable<> resume; Observable seeking; Observable seekedTo; Observable> &> timeRangeChanged; Observable eventsMerged; Observable *, bool> msgsReceived; Observable error; SourceSet sources; protected: void postToMainThread(std::function fn); void postToMainThreadAndWait(std::function fn); void cancelWaits(); void requestUpdateLastMessages() { postToMainThread([this]() { updateLastMessages(); }); } void mergeEvents(const std::vector &events); void insertEvents(const std::vector &events, const MessageEventsMap &msg_events); const CanEvent *newEvent(uint64_t mono_time, const cereal::CanData::Reader &c); void updateEvent(const MessageId &id, double sec, const uint8_t *data, uint8_t size); void waitForSeekFinshed(); virtual void updateLastMessages(); std::vector all_events_; double current_sec_ = 0; std::optional> time_range_; private: void updateLastMsgsTo(double sec); void updateMasks(); MessageEventsMap events_; std::unordered_map last_msgs; std::unique_ptr event_buffer_; std::shared_ptr alive_ = std::make_shared(true); Connections connections_; std::mutex mutex_; std::condition_variable wait_cv_; bool seek_finished_ = false; bool exiting_ = false; std::set new_msgs_; std::unordered_map messages_; std::unordered_map> masks_; }; class DummyStream : public AbstractStream { public: std::string routeName() const override { return "No Stream"; } void start() override {} }; extern AbstractStream *can;