From 7d7d50d55b9e0d65a708d43e84d6334584bb4720 Mon Sep 17 00:00:00 2001 From: Tomaz Cvetko Date: Wed, 17 Dec 2025 10:00:03 +0100 Subject: [PATCH 1/2] Rework unitree_device to include DDS listener --- .../unitree_module/sensor_data_listener.h | 50 +++ .../include/unitree_module/unitree_channel.h | 29 +- .../include/unitree_module/unitree_device.h | 66 +-- .../include/unitree_module/unitree_module.h | 1 + unitree_module/src/CMakeLists.txt | 6 + unitree_module/src/CycloneData.idl | 21 + unitree_module/src/sensor_data_listener.cpp | 80 ++++ unitree_module/src/unitree_channel.cpp | 72 ++- unitree_module/src/unitree_device.cpp | 412 ++++++------------ unitree_module/src/unitree_module.cpp | 23 +- 10 files changed, 358 insertions(+), 402 deletions(-) create mode 100644 unitree_module/include/unitree_module/sensor_data_listener.h create mode 100644 unitree_module/src/CycloneData.idl create mode 100644 unitree_module/src/sensor_data_listener.cpp diff --git a/unitree_module/include/unitree_module/sensor_data_listener.h b/unitree_module/include/unitree_module/sensor_data_listener.h new file mode 100644 index 0000000..2bb0023 --- /dev/null +++ b/unitree_module/include/unitree_module/sensor_data_listener.h @@ -0,0 +1,50 @@ +#include + +/* Include the C++ DDS API. */ +#include "dds/dds.hpp" +#include "dds/domain/qos/DomainParticipantQos.hpp" + +#include "CycloneData.hpp" + +BEGIN_NAMESPACE_UNITREE_MODULE + +namespace +{ + using MsgType_ = CycloneData::Msg; +} + +class SensorDataListener : public virtual dds::sub::DataReaderListener +{ +public: + // using MsgType = unitree_go::msg::dds_::LowState_; + using MsgType = MsgType_; + using ArgType = std::vector&&; + inline static std::string GetTopicName() + { + return "CycloneData_Msg"; + } + + explicit SensorDataListener(const std::function&&)>& func); + + void on_data_available(dds::sub::DataReader& reader) override; + + void on_subscription_matched(dds::sub::DataReader& reader, + const dds::core::status::SubscriptionMatchedStatus& status) override; + + void on_sample_lost(dds::sub::DataReader& reader, const dds::core::status::SampleLostStatus& status) override; + + void on_requested_deadline_missed(dds::sub::DataReader& reader, + const dds::core::status::RequestedDeadlineMissedStatus& status) override; + + void on_requested_incompatible_qos(dds::sub::DataReader& reader, + const dds::core::status::RequestedIncompatibleQosStatus& status) override; + + void on_sample_rejected(dds::sub::DataReader& reader, const dds::core::status::SampleRejectedStatus& status) override; + + void on_liveliness_changed(dds::sub::DataReader& reader, const dds::core::status::LivelinessChangedStatus& status) override; + +private: + std::function&&)> callback; +}; + +END_NAMESPACE_UNITREE_MODULE \ No newline at end of file diff --git a/unitree_module/include/unitree_module/unitree_channel.h b/unitree_module/include/unitree_module/unitree_channel.h index f11db85..13fa475 100644 --- a/unitree_module/include/unitree_module/unitree_channel.h +++ b/unitree_module/include/unitree_module/unitree_channel.h @@ -15,22 +15,23 @@ */ #pragma once -#include #include #include +#include #include #include +#include BEGIN_NAMESPACE_UNITREE_MODULE enum class WaveformType { Sine, Rect, None, Counter, ConstantValue }; -DECLARE_OPENDAQ_INTERFACE(IRefChannel, IBaseObject) +DECLARE_OPENDAQ_INTERFACE(IDogChannel, IBaseObject) { - virtual void collectSamples(std::chrono::microseconds curTime) = 0; + virtual void publishSamples(std::chrono::microseconds curTime, std::vector & data) = 0; }; -struct RefChannelInit +struct UnitreeChannelInit { size_t index; double sampleRate; @@ -38,16 +39,16 @@ struct RefChannelInit std::chrono::microseconds microSecondsFromEpochToStartTime; }; -class ExampleChannel final : public ChannelImpl +class UnitreeChannel final : public ChannelImpl { public: - explicit ExampleChannel(const ContextPtr& context, + explicit UnitreeChannel(const ContextPtr& context, const ComponentPtr& parent, const StringPtr& localId, - const RefChannelInit& init); + const UnitreeChannelInit& init); - // IRefChannel - void collectSamples(std::chrono::microseconds curTime) override; + // IDogChannel + void publishSamples(std::chrono::microseconds curTime, std::vector& data) override; static std::string getEpoch(); static RatioPtr getResolution(); @@ -72,13 +73,13 @@ class ExampleChannel final : public ChannelImpl uint64_t counter; double sampleRate; - SignalConfigPtr valueSignal; SignalConfigPtr timeSignal; + SignalConfigPtr valueSignal; - // SignalConfigPtr forceFLsignal; - // SignalConfigPtr forceFRsignal; - // SignalConfigPtr forceRLsignal; - // SignalConfigPtr forceRRsignal; + SignalConfigPtr forceFLsignal; + SignalConfigPtr forceFRsignal; + SignalConfigPtr forceRLsignal; + SignalConfigPtr forceRRsignal; }; END_NAMESPACE_UNITREE_MODULE diff --git a/unitree_module/include/unitree_module/unitree_device.h b/unitree_module/include/unitree_module/unitree_device.h index 532ea32..6241d0c 100644 --- a/unitree_module/include/unitree_module/unitree_device.h +++ b/unitree_module/include/unitree_module/unitree_device.h @@ -15,37 +15,22 @@ */ #pragma once -#include -#include -#include -#include #include +#include +#include -BEGIN_NAMESPACE_UNITREE_MODULE - -class myDDSDevice -{ -public: - explicit myDDSDevice(const std::function& function); - - void acqLoop(); +#include +#include - std::thread acqThread; - std::function function; - std::condition_variable cv; - std::mutex mutex; +#include - int16_t smpl; - int sign; -}; +BEGIN_NAMESPACE_UNITREE_MODULE class UnitreeDevice final : public Device { public: explicit UnitreeDevice(const ContextPtr& ctx, const ComponentPtr& parent); - // ~UnitreeDevice() override; - - // void UnitreeDevice_fun(const std::function& function); + ~UnitreeDevice() override; static DeviceInfoPtr CreateDeviceInfo(); static DeviceTypePtr CreateType(); @@ -54,41 +39,30 @@ class UnitreeDevice final : public Device DeviceInfoPtr onGetInfo() override; uint64_t onGetTicksSinceOrigin() override; - // std::condition_variable cv; - // std::mutex mutex; - - // int16_t smpl; - // int sign; - - std::shared_ptr myDev; - private: - // void initDomain(); - // void initChannels(); - void initSignals(); - // void acqLoop(); - void processData(int16_t data) const; + void initDomain(); + void initChannels(); + + // void processData(int16_t data) const; // void processData(std::vector data) const; + void acqLoop(); + void forwardDataCallback(std::vector&& data); + void forwardData(std::vector& data); + std::chrono::microseconds getMicroSecondsSinceDeviceStart() const; - std::thread acqThread; + bool stopAcq; + std::queue> queue; + + std::mutex mutex; std::condition_variable cv; + std::thread acqThread; std::chrono::steady_clock::time_point startTime; std::chrono::microseconds microSecondsFromEpochToDeviceStart; ChannelPtr channel1; - ChannelPtr channel2; - - SignalConfigPtr sig_F_fl; - SignalConfigPtr sig_F_fr; - SignalConfigPtr sig_F_rl; - SignalConfigPtr sig_F_rr; - SignalConfigPtr sigTime; - - size_t acqLoopTime; - bool stopAcq; }; END_NAMESPACE_UNITREE_MODULE diff --git a/unitree_module/include/unitree_module/unitree_module.h b/unitree_module/include/unitree_module/unitree_module.h index 49b5419..cb38434 100644 --- a/unitree_module/include/unitree_module/unitree_module.h +++ b/unitree_module/include/unitree_module/unitree_module.h @@ -32,6 +32,7 @@ class UnitreeModule final : public Module private: std::mutex sync; bool deviceAdded; + size_t deviceIndex; }; END_NAMESPACE_UNITREE_MODULE diff --git a/unitree_module/src/CMakeLists.txt b/unitree_module/src/CMakeLists.txt index aac9443..e48cc4c 100644 --- a/unitree_module/src/CMakeLists.txt +++ b/unitree_module/src/CMakeLists.txt @@ -1,17 +1,22 @@ set(LIB_NAME unitree_module) set(MODULE_HEADERS_DIR ../include/${TARGET_FOLDER_NAME}) +# Use Cyclone IDL compiler to generate source code +idlcxx_generate(TARGET cyclonedata FILES CycloneData.idl WARNINGS no-implicit-extensibility) + set(SRC_Include common.h module_dll.h unitree_module.h unitree_device.h unitree_channel.h + sensor_data_listener.h ) set(SRC_Srcs module_dll.cpp unitree_module.cpp unitree_device.cpp unitree_channel.cpp + sensor_data_listener.cpp ) prepend_include(${TARGET_FOLDER_NAME} SRC_Include) @@ -41,6 +46,7 @@ target_link_libraries(${LIB_NAME} PUBLIC daq::opendaq CycloneDDS-CXX::ddscxx + cyclonedata ) target_include_directories(${LIB_NAME} PUBLIC $ diff --git a/unitree_module/src/CycloneData.idl b/unitree_module/src/CycloneData.idl new file mode 100644 index 0000000..500900e --- /dev/null +++ b/unitree_module/src/CycloneData.idl @@ -0,0 +1,21 @@ +/* + * Copyright(c) 2006 to 2020 ZettaScale Technology and others + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License v. 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Eclipse Distribution License + * v. 1.0 which is available at + * http://www.eclipse.org/org/documents/edl-v10.php. + * + * SPDX-License-Identifier: EPL-2.0 OR BSD-3-Clause + */ +module CycloneData +{ + struct Msg + { + long userID; + string message; + short values[4]; + }; + #pragma keylist Msg userID +}; diff --git a/unitree_module/src/sensor_data_listener.cpp b/unitree_module/src/sensor_data_listener.cpp new file mode 100644 index 0000000..e325ba1 --- /dev/null +++ b/unitree_module/src/sensor_data_listener.cpp @@ -0,0 +1,80 @@ +#include + +BEGIN_NAMESPACE_UNITREE_MODULE + +SensorDataListener::SensorDataListener(const std::function&&)>& func) + : callback(func) +{ +} + +void SensorDataListener::on_data_available(dds::sub::DataReader& reader) +{ + dds::sub::LoanedSamples samples = reader.take(); + size_t numberOfSamples = samples.length(); + + constexpr size_t numberOfComponents = 4; + std::vector buffer(numberOfComponents * numberOfSamples, 0); + + size_t sampleIndex = 0; + for (const auto& sample : samples) + { + if (!sample.info().valid()) + { + continue; + } + + const MsgType& msg = sample.data(); + const std::array& values = msg.values(); + + // Pack into flat array as [x0, y0, z0, w0, x1, y1, z1, w1] + std::copy(values.data(), values.data() + 4, &buffer[4 * sampleIndex]); + + // Equivalent but a bit easer to control + // for (size_t i = 0; i < 4; ++i) + // { + // const size_t index = numberOfComponents * sampleIndex + i; + // buffer[i] - values[i]; + // } + + // Question: Is there any information about the time these values were measured that we must propagate? + + ++sampleIndex; + } + this->callback(std::move(buffer)); +} + +void SensorDataListener::on_subscription_matched(dds::sub::DataReader& reader, + const dds::core::status::SubscriptionMatchedStatus& status) +{ + std::cout << "on_subscription_matched" << std::endl; +} + +void SensorDataListener::on_sample_lost(dds::sub::DataReader& reader, const dds::core::status::SampleLostStatus& status) +{ + std::cout << "on_sample_lost" << std::endl; +} + +void SensorDataListener::on_requested_deadline_missed(dds::sub::DataReader& reader, + const dds::core::status::RequestedDeadlineMissedStatus& status) +{ + std::cout << "on_requested_deadline_missed" << std::endl; +} + +void SensorDataListener::on_requested_incompatible_qos(dds::sub::DataReader& reader, + const dds::core::status::RequestedIncompatibleQosStatus& status) +{ + std::cout << "on_requested_incompatible_qos" << std::endl; +} + +void SensorDataListener::on_sample_rejected(dds::sub::DataReader& reader, const dds::core::status::SampleRejectedStatus& status) +{ + std::cout << "on_sample_rejected" << std::endl; +} + +void SensorDataListener::on_liveliness_changed(dds::sub::DataReader& reader, + const dds::core::status::LivelinessChangedStatus& status) +{ + std::cout << "on_liveliness_changed" << std::endl; +} + +END_NAMESPACE_UNITREE_MODULE \ No newline at end of file diff --git a/unitree_module/src/unitree_channel.cpp b/unitree_module/src/unitree_channel.cpp index 4c2df26..e60131f 100644 --- a/unitree_module/src/unitree_channel.cpp +++ b/unitree_module/src/unitree_channel.cpp @@ -1,19 +1,19 @@ -#include -#include -#include -#include #include -#include #include +#include +#include +#include +#include +#include #include BEGIN_NAMESPACE_UNITREE_MODULE -ExampleChannel::ExampleChannel(const ContextPtr& context, +UnitreeChannel::UnitreeChannel(const ContextPtr& context, const ComponentPtr& parent, const StringPtr& localId, - const RefChannelInit& init) - : ChannelImpl(FunctionBlockType("ExampleChannel", fmt::format("AI{}", init.index + 1), ""), context, parent, localId) + const UnitreeChannelInit& init) + : ChannelImpl(FunctionBlockType("UnitreeChannel", fmt::format("AI{}", init.index + 1), ""), context, parent, localId) , index(init.index) , startTime(init.startTime) , microSecondsFromEpochToStartTime(init.microSecondsFromEpochToStartTime) @@ -26,37 +26,18 @@ ExampleChannel::ExampleChannel(const ContextPtr& context, buildSignalDescriptors(); } -void ExampleChannel::initProperties() +void UnitreeChannel::initProperties() { } - -uint64_t ExampleChannel::getSamplesSinceStart(std::chrono::microseconds time) const -{ - const uint64_t samplesSinceStart = static_cast(std::trunc(static_cast((time - startTime).count()) / 1000000.0 * sampleRate)); - return samplesSinceStart; -} - -void ExampleChannel::collectSamples(std::chrono::microseconds curTime) +void UnitreeChannel::publishSamples(std::chrono::microseconds curTime, std::vector& data) { auto lock = this->getAcquisitionLock(); - const uint64_t samplesSinceStart = getSamplesSinceStart(curTime); - auto newSamples = samplesSinceStart - samplesGenerated; - - if (newSamples > 0) - { - const auto packetTime = samplesGenerated * deltaT + static_cast(microSecondsFromEpochToStartTime.count()); - auto [dataPacket, domainPacket] = generateSamples(static_cast(packetTime), newSamples); - - valueSignal.sendPacket(std::move(dataPacket)); - timeSignal.sendPacket(std::move(domainPacket)); - - samplesGenerated = samplesSinceStart; - } + // TODO } -std::tuple ExampleChannel::generateSamples(int64_t curTime, uint64_t newSamples) +std::tuple UnitreeChannel::generateSamples(int64_t curTime, uint64_t newSamples) { auto domainPacket = DataPacket(timeSignal.getDescriptor(), newSamples, curTime); DataPacketPtr dataPacket = DataPacketWithDomain(domainPacket, valueSignal.getDescriptor(), newSamples); @@ -72,7 +53,7 @@ std::tuple ExampleChannel::generateSamples(int64_t curTime return {dataPacket, domainPacket}; } -std::string ExampleChannel::getEpoch() +std::string UnitreeChannel::getEpoch() { const std::time_t epochTime = std::chrono::system_clock::to_time_t(std::chrono::time_point{}); @@ -82,19 +63,19 @@ std::string ExampleChannel::getEpoch() return { buf }; } -RatioPtr ExampleChannel::getResolution() +RatioPtr UnitreeChannel::getResolution() { return Ratio(1, 1000000); } -Int ExampleChannel::getDeltaT(const double sr) const +Int UnitreeChannel::getDeltaT(const double sr) const { const double tickPeriod = getResolution(); const double samplePeriod = 1.0 / sr; return static_cast(std::round(samplePeriod / tickPeriod)); } -void ExampleChannel::buildSignalDescriptors() +void UnitreeChannel::buildSignalDescriptors() { const auto valueDescriptor = DataDescriptorBuilder().setSampleType(SampleType::Float64).setUnit(Unit("V", -1, "volts", "voltage")); @@ -124,20 +105,23 @@ void ExampleChannel::buildSignalDescriptors() timeSignal.setDescriptor(timeDescriptor.build()); } -void ExampleChannel::createSignals() +void UnitreeChannel::createSignals() { valueSignal = createAndAddSignal(fmt::format("AI{}", index)); timeSignal = createAndAddSignal(fmt::format("AI{}Time", index), nullptr, false); valueSignal.setDomainSignal(timeSignal); - // forceFLsignal = createAndAddSignal(fmt::format("FL force")); - // forceFRsignal = createAndAddSignal(fmt::format("FR force")); - // forceRLsignal = createAndAddSignal(fmt::format("RL force")); - // forceRRsignal = createAndAddSignal(fmt::format("RR force")); - // forceFLsignal.setDomainSignal(timeSignal); - // forceFRsignal.setDomainSignal(timeSignal); - // forceRLsignal.setDomainSignal(timeSignal); - // forceRRsignal.setDomainSignal(timeSignal); + forceFLsignal = createAndAddSignal(fmt::format("FL force")); + forceFLsignal.setDomainSignal(timeSignal); + + forceFRsignal = createAndAddSignal(fmt::format("FR force")); + forceFRsignal.setDomainSignal(timeSignal); + + forceRLsignal = createAndAddSignal(fmt::format("RL force")); + forceRLsignal.setDomainSignal(timeSignal); + + forceRRsignal = createAndAddSignal(fmt::format("RR force")); + forceRRsignal.setDomainSignal(timeSignal); } END_NAMESPACE_UNITREE_MODULE diff --git a/unitree_module/src/unitree_device.cpp b/unitree_module/src/unitree_device.cpp index 0a5c29b..9052532 100644 --- a/unitree_module/src/unitree_device.cpp +++ b/unitree_module/src/unitree_device.cpp @@ -1,15 +1,17 @@ -#include -#include #include -#include -#include -#include -#include -#include -#include #include +#include +#include +#include #include +#include #include + +#include +#include +#include + +#include #include // #include @@ -36,162 +38,48 @@ Get-ChildItem -Path D:\Razno\unitree_go2 -Filter CycloneDDSConfig.cmake -Recurse BEGIN_NAMESPACE_UNITREE_MODULE - - - -// -//class SensorDataListener : public virtual dds::sub::DataReaderListener -//{ -//public: -// virtual void on_data_available(dds::sub::DataReader& reader) -// { -// auto samples = reader.take(); -// for (const auto& sample : samples) -// { -// if (sample.info().valid()) -// { -// const auto& state = sample.data(); -// -// auto forces = state.foot_force(); -// std::cout << "foot forces:" -// << " FL=" << forces[0] << " FR=" << forces[1] << " RL=" << forces[2] << " RR=" << forces[3] << std::endl; -// -// // std::this_thread::sleep_for(std::chrono::milliseconds(200)); -// } -// } -// } -// -// virtual void on_subscription_matched(dds::sub::DataReader& reader, -// const dds::core::status::SubscriptionMatchedStatus& status) -// { -// std::cout << "on_subscription_matched" << std::endl; -// } -// -// virtual void on_sample_lost(dds::sub::DataReader& reader, -// const dds::core::status::SampleLostStatus& status) -// { -// std::cout << "on_sample_lost" << std::endl; -// } -// -// virtual void on_requested_deadline_missed(dds::sub::DataReader& reader, -// const dds::core::status::RequestedDeadlineMissedStatus& status) -// { -// std::cout << "on_requested_deadline_missed" << std::endl; -// } -// -// virtual void on_requested_incompatible_qos(dds::sub::DataReader& reader, -// const dds::core::status::RequestedIncompatibleQosStatus& status) -// { -// std::cout << "on_requested_incompatible_qos" << std::endl; -// } -// -// virtual void on_sample_rejected(dds::sub::DataReader& reader, -// const dds::core::status::SampleRejectedStatus& status) -// { -// std::cout << "on_sample_rejected" << std::endl; -// } -// -// virtual void on_liveliness_changed(dds::sub::DataReader& reader, -// const dds::core::status::LivelinessChangedStatus& status) -// { -// std::cout << "on_liveliness_changed" << std::endl; -// } -//}; -// - - - - - - -using namespace std::chrono_literals; -using milli = std::chrono::milliseconds; - -myDDSDevice::myDDSDevice(const std::function < void(int16_t)>& function) - : function(function) -{ - smpl = 0; - sign = 1; - acqThread = std::thread{&myDDSDevice::acqLoop, this}; -} - -static std::random_device rd; -static std::mt19937 gen(rd()); -static std::uniform_int_distribution dist(10, 30); - -void myDDSDevice::acqLoop() +namespace { - daqNameThread("RobotDog"); - - auto lock = std::unique_lock(mutex); - - auto loopTime = milli(10); - auto prevLoopTime = std::chrono::steady_clock::now(); - - while (true) - { - const auto time = std::chrono::steady_clock::now(); - cv.wait_until(lock, prevLoopTime + loopTime); - - if (smpl == 0) - sign = 1; - else if (smpl == 10) - sign = -1; - - smpl += sign; - function(smpl); - - prevLoopTime = time; - loopTime = milli(static_cast(dist(gen))); - } + constexpr bool CallbackQueueing = false; } UnitreeDevice::UnitreeDevice(const ContextPtr& ctx, const ComponentPtr& parent) : GenericDevice<>(ctx, parent, "DewesoftRobotics_B42D4000O4N9D402", nullptr, "RobotDog") - //, microSecondsFromEpochToDeviceStart(0) - //, acqLoopTime(50) - //, stopAcq(false) + , microSecondsFromEpochToDeviceStart(0) + , stopAcq(false) + , acqThread([this]() { this->acqLoop(); }) { - /*this->loggerComponent = ctx.getLogger().getOrAddComponent(UNITREE_MODULE_NAME); + this->loggerComponent = ctx.getLogger().getOrAddComponent(UNITREE_MODULE_NAME); initDomain(); initChannels(); - acqThread = std::thread{ &UnitreeDevice::acqLoop, this };*/ /****************************************/ // Unitree DDS - // dds::domain::DomainParticipant participant(0); - // dds::topic::Topic topic(participant, "rt/lowstate"); - // dds::sub::Subscriber subscriber(participant); - // dds::sub::DataReader reader( - // subscriber, topic, subscriber.default_datareader_qos(), new SensorDataListener(), dds::core::status::StatusMask::data_available()); - // std::cout << "Listening for Unitree Go2 sensor data on topic 'rt/lowstate'..." << std::endl; + dds::domain::DomainParticipant participant(0); + dds::topic::Topic topic(participant, SensorDataListener::GetTopicName()); + dds::sub::Subscriber subscriber(participant); + + dds::sub::DataReader reader( + subscriber, + topic, + subscriber.default_datareader_qos(), + new SensorDataListener([this](std::vector&& data) { this->forwardDataCallback(std::move(data)); }), + dds::core::status::StatusMask::data_available()); + std::cout << "Listening for Unitree Go2 sensor data on topic '" << SensorDataListener::GetTopicName() << "'..." << std::endl; /****************************************/ +} - myDev = std::make_shared([this](int16_t data) { processData(data); }); - initSignals(); - - // create epoch - const std::time_t epochTime = std::chrono::system_clock::to_time_t(std::chrono::time_point{}); - - char buf[48]; - strftime(buf, sizeof buf, "%Y-%m-%dT%H:%M:%SZ", gmtime(&epochTime)); - - std::string epoch = {buf}; +UnitreeDevice::~UnitreeDevice() +{ + // TODO: Proper shutdown + { + auto lock = this->getAcquisitionLock(); + stopAcq = true; + } - this->setDeviceDomain( - DeviceDomain(Ratio(1, 1'000'000), epoch, UnitBuilder().setName("seconds").setSymbol("s").setQuantity("time").build())); + acqThread.join(); } -//UnitreeDevice::~UnitreeDevice() -//{ -// { -// auto lock = this->getAcquisitionLock(); -// stopAcq = true; -// } -// -// acqThread.join(); -//} - DeviceInfoPtr UnitreeDevice::CreateDeviceInfo() { // auto devInfo = DeviceInfo("example://device"); @@ -209,72 +97,15 @@ DeviceInfoPtr UnitreeDevice::CreateDeviceInfo() DeviceTypePtr UnitreeDevice::CreateType() { - return DeviceType("unitree_dev", - "Unitree device", - "Unitree device", - "unitree"); + return DeviceType("robotDog", "Robot dog", "", "daq.dog"); } -//void UnitreeDevice::initChannels() -//{ -// auto microSecondsSinceDeviceStart = getMicroSecondsSinceDeviceStart(); -// -// RefChannelInit init{0, 1000, microSecondsSinceDeviceStart, microSecondsFromEpochToDeviceStart}; -// channel1 = createAndAddChannel(ioFolder, "ExCh1", init); -// -// init.index = 1; -// channel2 = createAndAddChannel(ioFolder, "ExCh2", init); -//} - -void UnitreeDevice::initSignals() +void UnitreeDevice::initChannels() { - sig_F_fl = Signal(context, this->signals, "Sig_F_fl"); - sig_F_fr = Signal(context, this->signals, "Sig_F_fr"); - sig_F_rl = Signal(context, this->signals, "Sig_F_rl"); - sig_F_rr = Signal(context, this->signals, "Sig_F_rr"); - sigTime = Signal(context, this->signals, "SigTime"); - - this->signals.addItem(sig_F_fl); - this->signals.addItem(sig_F_fr); - this->signals.addItem(sig_F_rl); - this->signals.addItem(sig_F_rr); - this->signals.addItem(sigTime); - - // Create resolution - auto resolution = Ratio(1, 1'000'000); - - // Create epoch - const std::time_t epochTime = std::chrono::system_clock::to_time_t(std::chrono::time_point{}); - - char buf[48]; - strftime(buf, sizeof buf, "%Y-%m-%dT%H:%M:%SZ", gmtime(&epochTime)); - - std::string epoch = {buf}; - - const auto timeDescriptor = DataDescriptorBuilder() - .setSampleType(SampleType::Int64) - .setUnit(Unit("s", -1, "seconds", "time")) - .setTickResolution(resolution) - .setRule(ExplicitDataRule()) - .setOrigin(epoch) - .build(); - - sigTime.setDescriptor(timeDescriptor); - - const auto forceDescriptor = DataDescriptorBuilder() - .setSampleType(SampleType::Int16) - .setUnit(Unit("N", -1, "Newton", "force")) - .build(); - - sig_F_fl.setDescriptor(forceDescriptor); - sig_F_fr.setDescriptor(forceDescriptor); - sig_F_rl.setDescriptor(forceDescriptor); - sig_F_rr.setDescriptor(forceDescriptor); - - sig_F_fl.setDomainSignal(sigTime); - sig_F_fr.setDomainSignal(sigTime); - sig_F_rl.setDomainSignal(sigTime); - sig_F_rr.setDomainSignal(sigTime); + auto microSecondsSinceDeviceStart = getMicroSecondsSinceDeviceStart(); + + UnitreeChannelInit init{0, 1000, microSecondsSinceDeviceStart, microSecondsFromEpochToDeviceStart}; + channel1 = createAndAddChannel(ioFolder, "ExCh1", init); } DeviceInfoPtr UnitreeDevice::onGetInfo() @@ -297,81 +128,94 @@ uint64_t UnitreeDevice::onGetTicksSinceOrigin() return duration_cast(system_clock::now().time_since_epoch()).count(); } -//void UnitreeDevice::initDomain() -//{ -// startTime = std::chrono::steady_clock::now(); -// auto startAbsTime = std::chrono::system_clock::now(); -// -// microSecondsFromEpochToDeviceStart = std::chrono::duration_cast(startAbsTime.time_since_epoch()); -// -// this->setDeviceDomain( -// DeviceDomain(ExampleChannel::getResolution(), -// ExampleChannel::getEpoch(), -// UnitBuilder().setName("second").setSymbol("s").setQuantity("time").build())); -//} - -//void UnitreeDevice::UnitreeDevice_fun(const std::function& function) -//{ -// : function(function) -// smpl = 0; -// sign = 1; -// acqThread = std::thread{&UnitreeDevice::acqLoop, this}; -//} - -//static std::random_device rd; // non-deterministic random source -//static std::mt19937 gen(rd()); // Mersenne Twister engine seeded once -//static std::uniform_int_distribution dist(10, 30); // inclusive range - -//void UnitreeDevice::acqLoop() -//{ -// const auto loopTime = milli(acqLoopTime); -// auto lock = getUniqueLock(); -// -// while (!stopAcq) -// { -// cv.wait_for(lock, loopTime); -// if (!stopAcq) -// { -// auto curTime = getMicroSecondsSinceDeviceStart(); -// channel1.asPtr()->collectSamples(curTime); -// channel2.asPtr()->collectSamples(curTime); -// } -// } -//} - -void UnitreeDevice::processData(int16_t data) const -// void UnitreeDevice::processData(std::vector data) const +void UnitreeDevice::initDomain() { - using namespace std::chrono; - int64_t time = duration_cast(system_clock::now().time_since_epoch()).count(); - - DataPacketPtr domainPacket = DataPacket(sigTime.getDescriptor(), 1); - int64_t* domainData = static_cast(domainPacket.getRawData()); - *domainData = time; - - DataPacketPtr dp_F_fl = DataPacketWithDomain(domainPacket, sig_F_fl.getDescriptor(), 1); - DataPacketPtr dp_F_fr = DataPacketWithDomain(domainPacket, sig_F_fr.getDescriptor(), 1); - DataPacketPtr dp_F_rl = DataPacketWithDomain(domainPacket, sig_F_rl.getDescriptor(), 1); - DataPacketPtr dp_F_rr = DataPacketWithDomain(domainPacket, sig_F_rr.getDescriptor(), 1); - - int16_t* data_F_fl = static_cast(dp_F_fl.getRawData()); - int16_t* data_F_fr = static_cast(dp_F_fr.getRawData()); - int16_t* data_F_rl = static_cast(dp_F_rl.getRawData()); - int16_t* data_F_rr = static_cast(dp_F_rr.getRawData()); - - // std::cout << "data: " << data_F_fl << " " << data_F_fr << " " << data_F_rl << " " << data_F_rr << std::endl; - - *data_F_fl = data; - *data_F_fr = data; - *data_F_rl = data; - *data_F_rr = data; - - sigTime.sendPacket(domainPacket); - sig_F_fl.sendPacket(dp_F_fl); - sig_F_fr.sendPacket(dp_F_fr); - sig_F_rl.sendPacket(dp_F_rl); - sig_F_rr.sendPacket(dp_F_rr); + startTime = std::chrono::steady_clock::now(); + auto startAbsTime = std::chrono::system_clock::now(); + + microSecondsFromEpochToDeviceStart = std::chrono::duration_cast(startAbsTime.time_since_epoch()); + + this->setDeviceDomain(DeviceDomain(UnitreeChannel::getResolution(), UnitreeChannel::getEpoch(), Unit("s", -1, "seconds", "time"))); +} + +void UnitreeDevice::acqLoop() +{ + daqNameThread("UnitreeDogThread"); + + while (true) + { + std::unique_lock lock{mutex}; + cv.wait(lock, [this]() { return !queue.empty() || stopAcq; }); + + if (stopAcq) + { + return; + } + + auto data = std::move(queue.front()); + queue.pop(); + lock.unlock(); + + forwardData(data); + } +} + +void UnitreeDevice::forwardDataCallback(std::vector&& data) +{ + if (CallbackQueueing) + { + std::unique_lock lock(mutex); + stopAcq = false; + queue.emplace(std::move(data)); + lock.unlock(); + cv.notify_all(); + return; + } + else + { + forwardData(data); + } +} + +void UnitreeDevice::forwardData(std::vector& data) +{ + // TODO: Maybe we need some multithreading locks here + auto curTime = getMicroSecondsSinceDeviceStart(); + channel1.asPtr()->publishSamples(curTime, data); } +// void UnitreeDevice::processData(int16_t data) const +// // void UnitreeDevice::processData(std::vector data) const +// { +// using namespace std::chrono; +// int64_t time = duration_cast(system_clock::now().time_since_epoch()).count(); + +// DataPacketPtr domainPacket = DataPacket(sigTime.getDescriptor(), 1); +// int64_t* domainData = static_cast(domainPacket.getRawData()); +// *domainData = time; + +// DataPacketPtr dp_F_fl = DataPacketWithDomain(domainPacket, sig_F_fl.getDescriptor(), 1); +// DataPacketPtr dp_F_fr = DataPacketWithDomain(domainPacket, sig_F_fr.getDescriptor(), 1); +// DataPacketPtr dp_F_rl = DataPacketWithDomain(domainPacket, sig_F_rl.getDescriptor(), 1); +// DataPacketPtr dp_F_rr = DataPacketWithDomain(domainPacket, sig_F_rr.getDescriptor(), 1); + +// int16_t* data_F_fl = static_cast(dp_F_fl.getRawData()); +// int16_t* data_F_fr = static_cast(dp_F_fr.getRawData()); +// int16_t* data_F_rl = static_cast(dp_F_rl.getRawData()); +// int16_t* data_F_rr = static_cast(dp_F_rr.getRawData()); + +// // std::cout << "data: " << data_F_fl << " " << data_F_fr << " " << data_F_rl << " " << data_F_rr << std::endl; + +// *data_F_fl = data; +// *data_F_fr = data; +// *data_F_rl = data; +// *data_F_rr = data; + +// sigTime.sendPacket(domainPacket); +// sig_F_fl.sendPacket(dp_F_fl); +// sig_F_fr.sendPacket(dp_F_fr); +// sig_F_rl.sendPacket(dp_F_rl); +// sig_F_rr.sendPacket(dp_F_rr); +// } END_NAMESPACE_UNITREE_MODULE diff --git a/unitree_module/src/unitree_module.cpp b/unitree_module/src/unitree_module.cpp index 71b09ce..84662cc 100644 --- a/unitree_module/src/unitree_module.cpp +++ b/unitree_module/src/unitree_module.cpp @@ -14,12 +14,16 @@ UnitreeModule::UnitreeModule(ContextPtr context) VersionInfo(UNITREE_MODULE_MAJOR_VERSION, UNITREE_MODULE_MINOR_VERSION, UNITREE_MODULE_PATCH_VERSION), std::move(context), UNITREE_MODULE_NAME) - //, deviceAdded(false) + , deviceAdded(false) + , deviceIndex(0) { } ListPtr UnitreeModule::onGetAvailableDevices() { + // TODO: Different modules handle this differently - potentially one device should be created immediately and + // always available? + /*ListPtr availableDevices = List(); std::cout << "Hello from UnitreeModule::onGetAvailableDevices()" << std::endl; @@ -32,14 +36,7 @@ ListPtr UnitreeModule::onGetAvailableDevices() DictPtr UnitreeModule::onGetAvailableDeviceTypes() { - /*auto result = Dict(); - - auto deviceType = UnitreeDevice::CreateType(); - result.set(deviceType.getId(), deviceType); - - return result;*/ - - auto type = DeviceType("robotDog", "Robot dog", "", "daq.dog"); + auto type = UnitreeDevice::CreateType(); return Dict({{type.getId(), type}}); } @@ -47,10 +44,10 @@ DevicePtr UnitreeModule::onCreateDevice(const StringPtr& connectionString, const ComponentPtr& parent, const PropertyObjectPtr& /*config*/) { - /*std::scoped_lock lock(sync); + std::scoped_lock lock(sync); std::string connStr = connectionString; - if (connStr.find("example://") != 0) + if (connStr.find("daq.dog://") != 0) throw std::runtime_error("Invalid connection string prefix"); if (deviceAdded) @@ -59,9 +56,7 @@ DevicePtr UnitreeModule::onCreateDevice(const StringPtr& connectionString, auto devicePtr = createWithImplementation(context, parent); deviceAdded = true; - return devicePtr;*/ - - return createWithImplementation(context, parent); + return devicePtr; } END_NAMESPACE_UNITREE_MODULE From 9684e575bf867ea929486923ef2e394aec318e38 Mon Sep 17 00:00:00 2001 From: Tomaz Cvetko Date: Wed, 17 Dec 2025 10:30:23 +0100 Subject: [PATCH 2/2] Theoretical packet creation. --- .../include/unitree_module/unitree_channel.h | 2 - unitree_module/src/unitree_channel.cpp | 75 +++++++++++-------- 2 files changed, 44 insertions(+), 33 deletions(-) diff --git a/unitree_module/include/unitree_module/unitree_channel.h b/unitree_module/include/unitree_module/unitree_channel.h index 13fa475..f6e08ec 100644 --- a/unitree_module/include/unitree_module/unitree_channel.h +++ b/unitree_module/include/unitree_module/unitree_channel.h @@ -59,7 +59,6 @@ class UnitreeChannel final : public ChannelImpl void buildSignalDescriptors(); uint64_t getSamplesSinceStart(std::chrono::microseconds time) const; - std::tuple generateSamples(int64_t curTime, uint64_t newSamples); [[nodiscard]] Int getDeltaT(const double sr) const; uint64_t deltaT; @@ -74,7 +73,6 @@ class UnitreeChannel final : public ChannelImpl double sampleRate; SignalConfigPtr timeSignal; - SignalConfigPtr valueSignal; SignalConfigPtr forceFLsignal; SignalConfigPtr forceFRsignal; diff --git a/unitree_module/src/unitree_channel.cpp b/unitree_module/src/unitree_channel.cpp index e60131f..39e9060 100644 --- a/unitree_module/src/unitree_channel.cpp +++ b/unitree_module/src/unitree_channel.cpp @@ -13,7 +13,7 @@ UnitreeChannel::UnitreeChannel(const ContextPtr& context, const ComponentPtr& parent, const StringPtr& localId, const UnitreeChannelInit& init) - : ChannelImpl(FunctionBlockType("UnitreeChannel", fmt::format("AI{}", init.index + 1), ""), context, parent, localId) + : ChannelImpl(FunctionBlockType("UnitreeChannel", fmt::format("DogCh{}", init.index + 1), ""), context, parent, localId) , index(init.index) , startTime(init.startTime) , microSecondsFromEpochToStartTime(init.microSecondsFromEpochToStartTime) @@ -34,23 +34,37 @@ void UnitreeChannel::publishSamples(std::chrono::microseconds curTime, std::vect { auto lock = this->getAcquisitionLock(); - // TODO -} - -std::tuple UnitreeChannel::generateSamples(int64_t curTime, uint64_t newSamples) -{ - auto domainPacket = DataPacket(timeSignal.getDescriptor(), newSamples, curTime); - DataPacketPtr dataPacket = DataPacketWithDomain(domainPacket, valueSignal.getDescriptor(), newSamples); - - double* buffer = static_cast(dataPacket.getRawData()); - - //for (uint64_t i = 0; i < newSamples; i++) - // buffer[i] = static_cast(counter++) / sampleRate; - - for (uint64_t i = 0; i < newSamples; i++) - buffer[i] = static_cast(5); - - return {dataPacket, domainPacket}; + const size_t bufferSize = data.size(); + assert(bufferSize % 4 == 0); + + const size_t sampleCount = bufferSize / 4; + + // TODO: I have no idea what timestamps are appropriate here + auto domainPacket = DataPacket(timeSignal.getDescriptor(), sampleCount, curTime.count()); + std::vector packets; + packets.reserve(4); + packets.push_back(DataPacketWithDomain(domainPacket, forceFLsignal.getDescriptor(), sampleCount)); // flPacket + packets.push_back(DataPacketWithDomain(domainPacket, forceFRsignal.getDescriptor(), sampleCount)); // frPacket + packets.push_back(DataPacketWithDomain(domainPacket, forceRLsignal.getDescriptor(), sampleCount)); // rlPacket + packets.push_back(DataPacketWithDomain(domainPacket, forceRRsignal.getDescriptor(), sampleCount)); // rrPacket + + // NOTE: If this is too slow, we can optimize the data structure + for (size_t sampleIndex = 0; sampleIndex < sampleCount; ++sampleIndex) + { + for (size_t i = 0; i < 4; ++i) + { + int16_t* packetData = static_cast(packets[i].getRawData()); + packetData[sampleIndex] = data[4 * sampleIndex + i]; + } + } + + timeSignal.sendPacket(domainPacket); + forceFLsignal.sendPacket(packets[0]); + forceFRsignal.sendPacket(packets[1]); + forceRLsignal.sendPacket(packets[2]); + forceRRsignal.sendPacket(packets[3]); + + counter += sampleCount; } std::string UnitreeChannel::getEpoch() @@ -68,6 +82,11 @@ RatioPtr UnitreeChannel::getResolution() return Ratio(1, 1000000); } +uint64_t UnitreeChannel::getSamplesSinceStart(std::chrono::microseconds time) const +{ + return counter; +} + Int UnitreeChannel::getDeltaT(const double sr) const { const double tickPeriod = getResolution(); @@ -77,19 +96,15 @@ Int UnitreeChannel::getDeltaT(const double sr) const void UnitreeChannel::buildSignalDescriptors() { - const auto valueDescriptor = DataDescriptorBuilder().setSampleType(SampleType::Float64).setUnit(Unit("V", -1, "volts", "voltage")); + const auto valueDescriptorForce = DataDescriptorBuilder().setSampleType(SampleType::Int16).setUnit(Unit("N", -1)); - // const auto valueDescriptorForce = DataDescriptorBuilder().setSampleType(SampleType::Int16).setUnit(Unit("N", -1)); - - valueSignal.setDescriptor(valueDescriptor.build()); - deltaT = getDeltaT(sampleRate); - - // forceFLsignal.setDescriptor(valueDescriptorForce.build()); - // forceFRsignal.setDescriptor(valueDescriptorForce.build()); - // forceRLsignal.setDescriptor(valueDescriptorForce.build()); - // forceRRsignal.setDescriptor(valueDescriptorForce.build()); + forceFLsignal.setDescriptor(valueDescriptorForce.build()); + forceFRsignal.setDescriptor(valueDescriptorForce.build()); + forceRLsignal.setDescriptor(valueDescriptorForce.build()); + forceRRsignal.setDescriptor(valueDescriptorForce.build()); // delta, start, tickResolution, unit, origin + deltaT = getDeltaT(sampleRate); // PacketOffset // PacketOffset + rule.delta * sampleIndex + rule.start -> ticks since origin @@ -107,9 +122,7 @@ void UnitreeChannel::buildSignalDescriptors() void UnitreeChannel::createSignals() { - valueSignal = createAndAddSignal(fmt::format("AI{}", index)); - timeSignal = createAndAddSignal(fmt::format("AI{}Time", index), nullptr, false); - valueSignal.setDomainSignal(timeSignal); + timeSignal = createAndAddSignal(fmt::format("Dog Time"), nullptr, false); forceFLsignal = createAndAddSignal(fmt::format("FL force")); forceFLsignal.setDomainSignal(timeSignal);