14#include "NetworkDataStreamInstrument.pb.h"
15#include "NetworkDataStreamInstrument.grpc.pb.h"
19 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
20 class NetworkDataStreamInstrumentT;
35 default:
throw Util::InvalidDataException(
"The given unit is not supported here. Did you forget to adjust this function or the IntensityUnitType enumeration in file \"Common.proto\"?");
39 namespace NetworkDataStreamInstrumentTasks
41 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
50 StubPtr = InstrData->template GetStub<DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument>();
53 auto Response =
InvokeStubFunc(StubPtr, &DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument::Stub::GetStreamInfo, {});
59 InstrData->RemoteStreamInfo.HardwareMinValue = Response.hardwareminvalue();
60 InstrData->RemoteStreamInfo.HardwareMaxValue = Response.hardwaremaxvalue();
61 InstrData->RemoteStreamInfo.IsBasicSampleTimeUsed = Response.isbasicsampletimeused();
62 InstrData->RemoteStreamInfo.StreamSizeRead = Util::NumToT<size_t>(Response.streamsizemsg().streamsizeread());
63 InstrData->RemoteStreamInfo.StreamSizeWrite = Util::NumToT<size_t>(Response.streamsizemsg().streamsizewrite());
65 InstrData->GetSampleStream()->SetStreamSize(InstrData->RemoteStreamInfo.StreamSizeWrite);
74 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
85 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
93 StubPtr = InstrData->template GetStub<DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument>();
96 auto StreamSizeResponse =
InvokeStubFunc(StubPtr, &DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument::Stub::GetStreamSize, {});
97 auto FinishedResponse =
InvokeStubFunc(StubPtr, &DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument::Stub::HasFinished, {});
98 auto RunningResponse =
InvokeStubFunc(StubPtr, &DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument::Stub::IsRunning, {});
101 auto InstrData = dynamic_InstrumentData_cast<
NetworkDataStreamInstrumentT<BaseInstr, 0, gRPCStubs...>>(Instance.InstrumentDataGetter());
103 InstrData->RemoteStreamInfo.StreamSizeRead = Util::NumToT<size_t>(StreamSizeResponse.streamsizeread());
104 InstrData->RemoteStreamInfo.StreamSizeWrite = Util::NumToT<size_t>(StreamSizeResponse.streamsizewrite());
108 if (InstrData->GetSampleStream()->GetStreamSizeWrite() != InstrData->RemoteStreamInfo.StreamSizeWrite)
110 InstrData->GetSampleStream()->SetStreamSize(InstrData->RemoteStreamInfo.StreamSizeWrite);
111 InstrData->SetLastReadRemoteSampleID(0);
121 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
130 DynExpProto::NetworkDataStreamInstrument::ReadMessage ReadMsg;
135 StubPtr = InstrData->template GetStub<DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument>();
136 ReadMsg.set_startsampleid(Util::NumToT<google::protobuf::uint64>(InstrData->GetLastReadRemoteSampleID()));
139 auto ReadResultMsg =
InvokeStubFunc(StubPtr, &DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument::Stub::Read, ReadMsg);
142 for (
decltype(ReadResultMsg.samples_size()) i = 0; i < ReadResultMsg.samples_size(); ++i)
143 InstrData->GetSampleStream()->WriteBasicSample({ ReadResultMsg.samples(i).value(), ReadResultMsg.samples(i).time() });
145 InstrData->SetLastReadRemoteSampleID(Util::NumToT<size_t>(ReadResultMsg.lastsampleid()));
151 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
161 std::vector<NetworkDataStreamInstrumentDataSampleStreamType::SampleType> Samples;
164 auto SampleStream = InstrData->template GetCastSampleStream<NetworkDataStreamInstrumentDataSampleStreamType>();
166 if (SampleStream->GetNumSamplesWritten() == InstrData->GetLastWrittenSampleID())
168 if (SampleStream->GetNumSamplesWritten() < InstrData->GetLastWrittenSampleID())
169 InstrData->SetLastWrittenSampleID(0);
171 StubPtr = InstrData->template GetStub<DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument>();
173 Samples = SampleStream->ReadRecentBasicSamples(InstrData->GetLastWrittenSampleID());
174 InstrData->SetLastWrittenSampleID(SampleStream->GetNumSamplesWritten());
177 DynExpProto::NetworkDataStreamInstrument::WriteMessage WriteMsg;
178 for (
const auto& Sample : Samples)
180 auto BasicSampleMsg = WriteMsg.add_samples();
181 BasicSampleMsg->set_value(Sample.Value);
182 BasicSampleMsg->set_time(Sample.Time);
185 InvokeStubFunc(StubPtr, &DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument::Stub::Write, WriteMsg);
191 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
203 StubPtr = InstrData->template GetStub<DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument>();
206 InvokeStubFunc(StubPtr, &DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument::Stub::ClearData, {});
212 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
224 StubPtr = InstrData->template GetStub<DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument>();
227 InvokeStubFunc(StubPtr, &DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument::Stub::Start, {});
233 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
245 StubPtr = InstrData->template GetStub<DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument>();
248 InvokeStubFunc(StubPtr, &DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument::Stub::Stop, {});
254 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
266 StubPtr = InstrData->template GetStub<DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument>();
269 InvokeStubFunc(StubPtr, &DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument::Stub::Restart, {});
275 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
285 DynExpProto::NetworkDataStreamInstrument::StreamSizeMessage StreamSizeMsg;
286 StreamSizeMsg.set_streamsizewrite(Util::NumToT<google::protobuf::uint64>(
StreamSizeInSamples));
291 StubPtr = InstrData->template GetStub<DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument>();
294 auto Response =
InvokeStubFunc(StubPtr, &DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument::Stub::SetStreamSize, StreamSizeMsg);
298 if (InstrData->GetSampleStream()->GetStreamSizeWrite() != Response.streamsizewrite())
300 InstrData->GetSampleStream()->
SetStreamSize(Response.streamsizewrite());
301 InstrData->SetLastReadRemoteSampleID(0);
311 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
323 StubPtr = InstrData->template GetStub<DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument>();
326 auto Response =
InvokeStubFunc(StubPtr, &DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument::Stub::ResetStreamSize, {});
330 if (InstrData->GetSampleStream()->GetStreamSizeWrite() != Response.streamsizewrite())
332 InstrData->GetSampleStream()->
SetStreamSize(Response.streamsizewrite());
333 InstrData->SetLastReadRemoteSampleID(0);
342 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
398 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
405 virtual const char*
GetParamClassTag() const noexcept
override {
return "NetworkDataStreamInstrumentParams"; }
414 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
434 template <
typename BaseInstr,
typename std::enable_if_t<std::is_base_of_v<DataStreamInstrument, BaseInstr>,
int>,
typename... gRPCStubs>
442 constexpr static auto Name() noexcept {
return "Network Data Stream Instrument"; }
445 :
gRPCInstrument<BaseInstr, 0, gRPCStubs...>(OwnerThreadID, std::move(Params)) {}
454 virtual std::chrono::milliseconds
GetTaskQueueDelay()
const {
return std::chrono::milliseconds(500); }
458 auto InstrData = dynamic_InstrumentData_cast<NetworkDataStreamInstrumentT>(this->GetInstrumentData());
459 return InstrData->HasFinished();
464 auto InstrData = dynamic_InstrumentData_cast<NetworkDataStreamInstrumentT>(this->GetInstrumentData());
465 return InstrData->IsRunning();
497 using StubType = DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument;
506 auto InstrData = dynamic_InstrumentData_cast<NetworkDataStreamInstrumentT>(this->GetInstrumentData());
507 return InstrData->GetValueUnit();
Implementation of a data stream meta instrument and of data streams input/output devices might work o...
Implements a circular data stream based on Util::circularbuf using samples of type BasicSample.
Implementation of the data stream meta instrument, which is a base class for all instruments reading/...
virtual DynExp::ParamsBasePtrType MakeParams(DynExp::ItemIDType ID, const DynExp::DynExpCore &Core) const override
Override to make derived classes call DynExp::MakeParams with the correct configurator type derived f...
virtual ~NetworkDataStreamInstrumentConfigurator()=default
NetworkDataStreamInstrumentConfigurator()=default
Util::OptionalBool Finished
NetworkDataStreamInstrumentData(size_t BufferSizeInSamples=1)
void ResetImpl(DynExp::InstrumentDataBase::dispatch_tag< gRPCInstrumentData< BaseInstr, 0, gRPCStubs... > >) override final
size_t LastReadRemoteSampleID
ID of the last sample read from the remote site and written to the assigned data stream.
RemoteStreamInfoType RemoteStreamInfo
size_t LastWrittenSampleID
ID of the last sample read from the assigned data stream and written to the remote site.
void SetLastWrittenSampleID(size_t SampleID) noexcept
virtual void ResetImpl(DynExp::InstrumentDataBase::dispatch_tag< NetworkDataStreamInstrumentData >)
const auto & GetRemoteStreamInfo() const noexcept
virtual ~NetworkDataStreamInstrumentData()=default
DynExp::Units::UnitType ValueUnit
auto HasFinished() const noexcept
void SetLastReadRemoteSampleID(size_t SampleID) noexcept
Util::OptionalBool Running
auto GetLastReadRemoteSampleID() const noexcept
bool IsBasicSampleTimeUsed
auto IsRunning() const noexcept
auto GetLastWrittenSampleID() const noexcept
DynExp::ParamsBase::DummyParam Dummy
virtual const char * GetParamClassTag() const noexcept override
This function is intended to be overridden once in each derived class returning the name of the respe...
virtual ~NetworkDataStreamInstrumentParams()=default
NetworkDataStreamInstrumentParams(DynExp::ItemIDType ID, const DynExp::DynExpCore &Core)
void ConfigureParamsImpl(DynExp::ParamsBase::dispatch_tag< gRPCInstrumentParams< BaseInstr, 0, gRPCStubs... > >) override final
virtual void ConfigureParamsImpl(DynExp::ParamsBase::dispatch_tag< NetworkDataStreamInstrumentParams >)
Data stream instrument for bidirectional gRPC communication.
static constexpr auto Name() noexcept
virtual std::unique_ptr< DynExp::ExitTaskBase > MakeExitTask() const override
Factory function for an exit task (ExitTaskBase). Override to define the desired deinitialization tas...
virtual std::unique_ptr< DynExp::UpdateTaskBase > MakeUpdateTask() const override
Factory function for an update task (UpdateTaskBase). Override to define the desired update task in d...
virtual std::string GetName() const override
Returns the name of this Object type.
NetworkDataStreamInstrumentT(const std::thread::id OwnerThreadID, DynExp::ParamsBasePtrType &&Params)
virtual void ClearData(DynExp::TaskBase::CallbackType CallbackFunc=nullptr) const override
virtual void Stop(DynExp::TaskBase::CallbackType CallbackFunc=nullptr) const override
virtual Util::OptionalBool IsRunning() const override
virtual ~NetworkDataStreamInstrumentT()
virtual Util::OptionalBool HasFinished() const override
virtual void Start(DynExp::TaskBase::CallbackType CallbackFunc=nullptr) const override
virtual void SetStreamSize(size_t BufferSizeInSamples, DynExp::TaskBase::CallbackType CallbackFunc=nullptr) const override
virtual std::unique_ptr< DynExp::InitTaskBase > MakeInitTask() const override
Factory function for an init task (InitTaskBase). Override to define the desired initialization task ...
virtual void ResetStreamSize(DynExp::TaskBase::CallbackType CallbackFunc=nullptr) const override
virtual void ResetImpl(DynExp::Object::dispatch_tag< NetworkDataStreamInstrumentT >)
void ResetImpl(DynExp::Object::dispatch_tag< gRPCInstrument< BaseInstr, 0, gRPCStubs... > >) override final
virtual std::chrono::milliseconds GetTaskQueueDelay() const
Read remote instrument's state periodically.
virtual void WriteData(DynExp::TaskBase::CallbackType CallbackFunc=nullptr) const override
virtual void Restart(DynExp::TaskBase::CallbackType CallbackFunc=nullptr) const override
virtual void ReadData(DynExp::TaskBase::CallbackType CallbackFunc=nullptr) const override
virtual DynExp::TaskResultType RunChild(DynExp::InstrumentInstance &Instance) override
Runs the task. Override RunChild() to define a derived task's action(s). Any exception leaving RunChi...
ClearTask(CallbackType CallbackFunc) noexcept
virtual void ExitFuncImpl(DynExp::ExitTaskBase::dispatch_tag< ExitTask >, DynExp::InstrumentInstance &Instance)
void ExitFuncImpl(DynExp::ExitTaskBase::dispatch_tag< gRPCInstrumentTasks::ExitTask< BaseInstr, 0, gRPCStubs... > >, DynExp::InstrumentInstance &Instance) override final
virtual void InitFuncImpl(DynExp::InitTaskBase::dispatch_tag< InitTask >, DynExp::InstrumentInstance &Instance)
void InitFuncImpl(DynExp::InitTaskBase::dispatch_tag< gRPCInstrumentTasks::InitTask< BaseInstr, 0, gRPCStubs... > >, DynExp::InstrumentInstance &Instance) override final
ReadTask(CallbackType CallbackFunc) noexcept
virtual DynExp::TaskResultType RunChild(DynExp::InstrumentInstance &Instance) override
Runs the task. Override RunChild() to define a derived task's action(s). Any exception leaving RunChi...
virtual DynExp::TaskResultType RunChild(DynExp::InstrumentInstance &Instance) override
Runs the task. Override RunChild() to define a derived task's action(s). Any exception leaving RunChi...
ResetStreamSizeTask(CallbackType CallbackFunc) noexcept
RestartTask(CallbackType CallbackFunc) noexcept
virtual DynExp::TaskResultType RunChild(DynExp::InstrumentInstance &Instance) override
Runs the task. Override RunChild() to define a derived task's action(s). Any exception leaving RunChi...
const size_t StreamSizeInSamples
virtual DynExp::TaskResultType RunChild(DynExp::InstrumentInstance &Instance) override
Runs the task. Override RunChild() to define a derived task's action(s). Any exception leaving RunChi...
SetStreamSizeTask(size_t StreamSizeInSamples, CallbackType CallbackFunc) noexcept
virtual DynExp::TaskResultType RunChild(DynExp::InstrumentInstance &Instance) override
Runs the task. Override RunChild() to define a derived task's action(s). Any exception leaving RunChi...
StartTask(CallbackType CallbackFunc) noexcept
virtual DynExp::TaskResultType RunChild(DynExp::InstrumentInstance &Instance) override
Runs the task. Override RunChild() to define a derived task's action(s). Any exception leaving RunChi...
StopTask(CallbackType CallbackFunc) noexcept
void UpdateFuncImpl(DynExp::UpdateTaskBase::dispatch_tag< gRPCInstrumentTasks::UpdateTask< BaseInstr, 0, gRPCStubs... > >, DynExp::InstrumentInstance &Instance) override final
virtual void UpdateFuncImpl(DynExp::UpdateTaskBase::dispatch_tag< UpdateTask >, DynExp::InstrumentInstance &Instance)
virtual DynExp::TaskResultType RunChild(DynExp::InstrumentInstance &Instance) override
Runs the task. Override RunChild() to define a derived task's action(s). Any exception leaving RunChi...
WriteTask(CallbackType CallbackFunc) noexcept
Explicit instantiation of derivable class NetworkDataStreamInstrumentT to create the network data str...
DynExpProto::NetworkDataStreamInstrument::NetworkDataStreamInstrument StubType
virtual ~NetworkDataStreamInstrument()
virtual DynExp::Units::UnitType GetValueUnit() const override
NetworkDataStreamInstrument(const std::thread::id OwnerThreadID, DynExp::ParamsBasePtrType &&Params)
Configurator class for gRPCInstrument.
Data class for gRPCInstrument.
Parameter class for gRPCInstrument.
Defines a task for deinitializing an instrument within an instrument inheritance hierarchy....
Defines a task for initializing an instrument within an instrument inheritance hierarchy....
Defines a task for updating an instrument within an instrument inheritance hierarchy....
Meta instrument template for transforming meta instruments into network instruments,...
DynExp's core class acts as the interface between the user interface and DynExp's internal data like ...
Refer to DynExp::ParamsBase::dispatch_tag.
Refer to DynExp::ParamsBase::dispatch_tag.
void MakeAndEnqueueTask(ArgTs &&...Args) const
Calls MakeTask() to construct a new task and subsequently enqueues the task into the instrument's tas...
Refer to ParamsBase::dispatch_tag.
Defines data for a thread belonging to a InstrumentBase instance. Refer to RunnableInstance.
const InstrumentBase::InstrumentDataGetterType InstrumentDataGetter
Getter for instrument's data. Refer to InstrumentBase::InstrumentDataGetterType.
Refer to ParamsBase::dispatch_tag.
Dummy parameter which is to be owned once by parameter classes that do not contain any other paramete...
Tag for function dispatching mechanism within this class used when derived classes are not intended t...
Type owning a callback function which is invoked when a task has finished, failed,...
Base class for all tasks being processed by instruments. The class must not contain public virtual fu...
CallbackType CallbackFunc
This callback function is called after the task has finished (either successfully or not) with a poin...
TaskBase(CallbackType CallbackFunc=nullptr, std::chrono::system_clock::time_point DeferUntil={}) noexcept
Constructs an instrument task.
Defines the return type of task functions.
Refer to DynExp::ParamsBase::dispatch_tag.
Data to operate on is invalid for a specific purpose. This indicates a corrupted data structure or fu...
Data type which stores an optional bool value (unknown, false, true). The type evaluates to bool whil...
Defines a meta instrument template for transforming meta instruments into network instruments,...
DynExp's instrument namespace contains the implementation of DynExp instruments which extend DynExp's...
constexpr DynExp::Units::UnitType ToDataStreamInstrumentUnitType(DynExpProto::Common::IntensityUnitType Unit)
BasicSampleStream NetworkDataStreamInstrumentDataSampleStreamType
ResponseMsgType InvokeStubFunc(StubPtrType< gRPCStub > StubPtr, StubFuncPtrType< gRPCStub, RequestMsgType, ResponseMsgType > StubFunc, const RequestMsgType &RequestMsg)
Invokes a gRPC stub function as a remote procedure call. Waits for a fixed amount of time (2 seconds)...
std::shared_ptr< typename gRPCStub::Stub > StubPtrType
Alias for a pointer to a gRPC stub.
UnitType
Units which can be used with DynExp instruments.
@ Ampere
Electric current in Ampere (A)
@ Arbitrary
Arbitrary units (a.u.)
@ LogicLevel
Logic level (TTL) units (1 or 0)
@ Power_W
Power in Watt (W)
@ Volt
Voltage in Volt (V)
@ Counts
Count rate in counts per second (cps)
std::unique_ptr< ParamsBase > ParamsBasePtrType
Alias for a pointer to the parameter system base class ParamsBase.
size_t ItemIDType
ID type of objects/items managed by DynExp.
std::unique_ptr< TaskT > MakeTask(ArgTs &&...Args)
Factory function to create a task to be enqueued in an instrument's task queue.
Accumulates include statements to provide a precompiled header.