init
This commit is contained in:
+220
@@ -0,0 +1,220 @@
|
||||
#include "ovasCPluginExternalStimulations.h"
|
||||
|
||||
#include <boost/interprocess/ipc/message_queue.hpp>
|
||||
|
||||
#include <vector>
|
||||
#include <ctime>
|
||||
|
||||
#include <system/ovCTime.h>
|
||||
|
||||
#include "../ovasCSettingsHelper.h"
|
||||
#include "../ovasCSettingsHelperOperators.h"
|
||||
|
||||
namespace OpenViBE {
|
||||
namespace AcquisitionServer {
|
||||
namespace Plugins {
|
||||
|
||||
CPluginExternalStimulations::CPluginExternalStimulations(const Kernel::IKernelContext& ctx)
|
||||
: IAcquisitionServerPlugin(ctx, CString("AcquisitionServer_Plugin_ExternalStimulations")), m_ExternalStimulationsQueueName("openvibeExternalStimulations")
|
||||
{
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Info << "Loading plugin: ExternalStimulations (deprecated)\n";
|
||||
|
||||
m_settings.add("EnableExternalStimulations", &m_IsExternalStimulationsEnabled);
|
||||
m_settings.add("ExternalStimulationQueueName", &m_ExternalStimulationsQueueName);
|
||||
m_settings.load();
|
||||
}
|
||||
|
||||
// Hooks
|
||||
|
||||
|
||||
bool CPluginExternalStimulations::startHook(const std::vector<CString>& /*selectedChannelNames*/, const size_t /*sampling*/, const size_t /*nChannel*/, const size_t /*nSamplePerSentBlock*/)
|
||||
{
|
||||
if (m_IsExternalStimulationsEnabled)
|
||||
{
|
||||
ftime(&m_CTStartTime);
|
||||
m_IsESThreadRunning = true;
|
||||
m_ESthreadPtr.reset(new std::thread(std::bind(&CPluginExternalStimulations::readExternalStimulations, this)));
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Info << "External stimulations (deprecated) activated...\n";
|
||||
}
|
||||
m_ExternalStimulations.clear();
|
||||
|
||||
m_DebugExternalStimulationsSent = 0;
|
||||
m_DebugCurrentReadIPCStimulations = 0;
|
||||
m_DebugStimulationsLost = 0;
|
||||
m_DebugStimulationsReceivedEarlier = 0;
|
||||
m_DebugStimulationsReceivedLate = 0;
|
||||
m_DebugStimulationsReceivedWrongSize = 0;
|
||||
m_DebugStimulationsBuffered = 0;
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
void CPluginExternalStimulations::loopHook(std::deque<std::vector<float>>& /* vPendingBuffer */, CStimulationSet& stimulationSet, const uint64_t start, const uint64_t end, const uint64_t /* sampleTime */)
|
||||
{
|
||||
if (m_IsExternalStimulationsEnabled)
|
||||
{
|
||||
//m_kernelCtx.getLogManager() << Kernel::LogLevel_Error << "Checking for external stimulations:" << p << "\n";
|
||||
addExternalStimulations(&stimulationSet, m_kernelCtx.getLogManager(), start, end);
|
||||
}
|
||||
}
|
||||
|
||||
void CPluginExternalStimulations::stopHook()
|
||||
{
|
||||
if (m_IsExternalStimulationsEnabled)
|
||||
{
|
||||
m_IsESThreadRunning = false;
|
||||
if (m_ESthreadPtr) { m_ESthreadPtr->join(); }
|
||||
else { m_kernelCtx.getLogManager() << Kernel::LogLevel_Warning << "Warning: External Stims plugin stopHook() tried to join a NULL thread\n"; }
|
||||
}
|
||||
|
||||
//software tagging diagnosting
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Debug << " Total external ones received through IPC: " << m_DebugCurrentReadIPCStimulations << "\n";
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Debug << " Sent to Designer: " << m_DebugExternalStimulationsSent << "\n";
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Debug << " Lost because of invalid timestamp: " << m_DebugStimulationsLost << "\n";
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Debug << " Stimulations that came earlier: " << m_DebugStimulationsReceivedEarlier << "\n";
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Debug << " Stimulations that came later: " << m_DebugStimulationsReceivedLate << "\n";
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Debug << " Stimulations that had wrong size: " << m_DebugStimulationsReceivedWrongSize << "\n";
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Debug << " Buffered: " << m_DebugStimulationsBuffered << "\n";
|
||||
//end software tagging diagnosting
|
||||
}
|
||||
|
||||
// Plugin specific methods
|
||||
|
||||
void CPluginExternalStimulations::readExternalStimulations()
|
||||
{
|
||||
using namespace boost::interprocess;
|
||||
|
||||
//std::cout << "Creating External Stimulations thread" << std::endl;
|
||||
//std::cout << "Queue Name : " << m_ExternalStimulationsQueueName << std::endl;
|
||||
//char mq_name[255];
|
||||
//std::strcpy(mq_name, m_ExternalStimulationsQueueName.toASCIIString());
|
||||
const int chunkLength = 3;
|
||||
const int pauseTime = 5;
|
||||
|
||||
uint32_t priority;
|
||||
size_t recvdSize;
|
||||
|
||||
uint64_t chunk[chunkLength];
|
||||
|
||||
while (m_IsESThreadRunning)
|
||||
{
|
||||
bool success;
|
||||
try
|
||||
{
|
||||
//Open a message queue.
|
||||
message_queue mq(open_only //only open
|
||||
, m_ExternalStimulationsQueueName.toASCIIString() //name
|
||||
//,mq_name //name
|
||||
);
|
||||
|
||||
success = mq.try_receive(&chunk, sizeof(chunk), recvdSize, priority);
|
||||
}
|
||||
catch (interprocess_exception& /* ex */)
|
||||
{
|
||||
//m_IsESThreadRunning = false;
|
||||
//m_kernelCtx.getLogManager() << Kernel::LogLevel_Error << "Problem with message queue in external stimulations:" << ex.what() << "\n";
|
||||
System::Time::sleep(pauseTime);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (!success)
|
||||
{
|
||||
System::Time::sleep(pauseTime);
|
||||
continue;
|
||||
}
|
||||
|
||||
m_DebugCurrentReadIPCStimulations++;
|
||||
|
||||
if (recvdSize != sizeof(chunk))
|
||||
{
|
||||
//m_kernelCtx.getLogManager() << Kernel::LogLevel_Error << "Problem with type of received data when reqding external stimulation!\n";
|
||||
m_DebugStimulationsReceivedWrongSize++;
|
||||
}
|
||||
else
|
||||
{
|
||||
//m_kernelCtx.getLogManager() << Kernel::LogLevel_Warning << "received\n";
|
||||
|
||||
SExternalStimulation stim;
|
||||
|
||||
stim.identifier = chunk[1];
|
||||
const uint64_t receivedTime = chunk[2];
|
||||
|
||||
//1. calculate time
|
||||
const uint64_t ctStartTimeMs = (m_CTStartTime.time * 1000 + m_CTStartTime.millitm);
|
||||
|
||||
const int64_t timeTest = receivedTime - ctStartTimeMs;
|
||||
|
||||
if (timeTest < 0)
|
||||
{
|
||||
m_DebugStimulationsLost++;
|
||||
//m_kernelCtx.getLogManager() << Kernel::LogLevel_Warning << "AS: external stimulation time is invalid, probably stimulation is before reference point, total invalid so far: " << m_FlashesLost << "\n";
|
||||
System::Time::sleep(pauseTime);
|
||||
continue; //we skip this stimulation
|
||||
}
|
||||
//2. Convert to OpenVibe time
|
||||
const uint64_t ctEventTime = receivedTime - ctStartTimeMs;
|
||||
|
||||
const double time = double(ctEventTime) / double(1000);
|
||||
|
||||
const uint64_t ovTime = CTime(time).time();
|
||||
stim.timestamp = ovTime;
|
||||
|
||||
//3. Store, the main thread will process it
|
||||
{
|
||||
//lock
|
||||
std::lock_guard<std::mutex> lock(m_es_mutex);
|
||||
|
||||
m_ExternalStimulations.push_back(stim);
|
||||
m_DebugStimulationsBuffered++;
|
||||
m_esAvailable.notify_one();
|
||||
//unlock
|
||||
}
|
||||
|
||||
System::Time::sleep(pauseTime);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
void CPluginExternalStimulations::addExternalStimulations(CStimulationSet* ss, Kernel::ILogManager& /*logm*/, const uint64_t start, const uint64_t /*end*/)
|
||||
{
|
||||
const uint64_t durationMs = 40;
|
||||
{
|
||||
//lock
|
||||
std::lock_guard<std::mutex> lock(m_es_mutex);
|
||||
|
||||
for (auto i = m_ExternalStimulations.begin(); i != m_ExternalStimulations.end(); ++i)
|
||||
{
|
||||
// if time is current or any time in the future - send it (AS will buffer it)
|
||||
if (i->timestamp >= start)
|
||||
{
|
||||
//flashes_in_this_time_chunk++;
|
||||
//logm << Kernel::LogLevel_Error << "Stimulation added." << "\n";
|
||||
ss->appendStimulation(i->identifier, i->timestamp, durationMs);
|
||||
}
|
||||
else
|
||||
{
|
||||
//the stimulation is coming too late - after the current block being processed
|
||||
//we correct the timestamp to the current block and we send it
|
||||
m_DebugStimulationsReceivedLate++;
|
||||
ss->appendStimulation(i->identifier, start, durationMs);
|
||||
}
|
||||
m_DebugExternalStimulationsSent++;
|
||||
}
|
||||
|
||||
// Since we processed all stimulations, we can clear the queue
|
||||
m_ExternalStimulations.clear();
|
||||
|
||||
m_esAvailable.notify_one();
|
||||
//unlock
|
||||
}
|
||||
}
|
||||
|
||||
bool CPluginExternalStimulations::setExternalStimulationsEnabled(const bool active)
|
||||
{
|
||||
m_IsExternalStimulationsEnabled = active;
|
||||
return true;
|
||||
}
|
||||
|
||||
} // namespace Plugins
|
||||
} // namespace AcquisitionServer
|
||||
} // namespace OpenViBE
|
||||
+81
@@ -0,0 +1,81 @@
|
||||
#pragma once
|
||||
|
||||
/**
|
||||
* \brief Acquisition Server plugin adding the capability to receive stimulations from external sources
|
||||
*
|
||||
* \author Anton Andreev
|
||||
* \author Jozef Legeny
|
||||
*
|
||||
* \note This plugin is deprecated. The users are recommended to use the TCP Tagging plugin instead. (11.05.2016)
|
||||
*
|
||||
*/
|
||||
|
||||
#include <thread>
|
||||
#include <mutex>
|
||||
#include <condition_variable>
|
||||
|
||||
#include <sys/timeb.h>
|
||||
|
||||
#include "ovasIAcquisitionServerPlugin.h"
|
||||
|
||||
namespace OpenViBE {
|
||||
namespace AcquisitionServer {
|
||||
class CAcquisitionServer;
|
||||
|
||||
namespace Plugins {
|
||||
class CPluginExternalStimulations final : public IAcquisitionServerPlugin
|
||||
{
|
||||
// Plugin interface
|
||||
public:
|
||||
explicit CPluginExternalStimulations(const Kernel::IKernelContext& ctx);
|
||||
~CPluginExternalStimulations() override {}
|
||||
|
||||
bool startHook(const std::vector<CString>& selectedChannelNames, const size_t sampling, const size_t nChannel, const size_t nSamplePerSentBlock) override;
|
||||
void stopHook() override;
|
||||
void loopHook(std::deque<std::vector<float>>& vPendingBuffer, CStimulationSet& stimulationSet, const uint64_t start, const uint64_t end,
|
||||
const uint64_t sampleTime) override;
|
||||
void acceptNewConnectionHook() override { m_ExternalStimulations.clear(); }
|
||||
|
||||
|
||||
// Plugin implementation
|
||||
|
||||
|
||||
struct SExternalStimulation
|
||||
{
|
||||
uint64_t timestamp;
|
||||
uint64_t identifier;
|
||||
};
|
||||
|
||||
void addExternalStimulations(CStimulationSet* ss, Kernel::ILogManager& logm, const uint64_t start, const uint64_t end);
|
||||
void readExternalStimulations();
|
||||
|
||||
//void acquireExternalStimulationsVRPN(CStimulationSet* ss, Kernel::ILogManager& logm, uint64_t start, uint64_t end);
|
||||
|
||||
struct timeb m_CTStartTime; //time when the acquisition process started in local computer time
|
||||
|
||||
std::vector<SExternalStimulation> m_ExternalStimulations;
|
||||
|
||||
bool m_IsExternalStimulationsEnabled = false;
|
||||
CString m_ExternalStimulationsQueueName;
|
||||
|
||||
bool setExternalStimulationsEnabled(bool active);
|
||||
bool isExternalStimulationsEnabled() const { return m_IsExternalStimulationsEnabled; }
|
||||
|
||||
// Debugging of external stimulations
|
||||
int m_DebugStimulationsLost = 0;
|
||||
int m_DebugExternalStimulationsSent = 0;
|
||||
int m_DebugCurrentReadIPCStimulations = 0;
|
||||
int m_DebugStimulationsReceivedEarlier = 0;
|
||||
int m_DebugStimulationsReceivedLate = 0;
|
||||
int m_DebugStimulationsReceivedWrongSize = 0;
|
||||
int m_DebugStimulationsBuffered = 0;
|
||||
|
||||
//added for acquiring external stimulations
|
||||
std::unique_ptr<std::thread> m_ESthreadPtr;
|
||||
bool m_IsESThreadRunning = false;
|
||||
std::mutex m_es_mutex;
|
||||
std::condition_variable m_esAvailable;
|
||||
};
|
||||
} // namespace Plugins
|
||||
} // namespace AcquisitionServer
|
||||
} // namespace OpenViBE
|
||||
+52
@@ -0,0 +1,52 @@
|
||||
#!/usr/bin/env python
|
||||
|
||||
# Example of tcp tagging client
|
||||
# The tag format is the same as with Shared Memory Tagging. It comprises three blocks of 8 bytes:
|
||||
#
|
||||
# ----------------------------------------------------------------------
|
||||
# | padding (8 bytes) | event id (8 bytes) | timestamp (8 bytes) |
|
||||
# ----------------------------------------------------------------------
|
||||
#
|
||||
# The padding is only for consistency with Shared Memory Tagging and has no utility.
|
||||
# The event id informs about the type of event happening.
|
||||
# The timestamp is the posix time (ms since Epoch) at the moment of the event.
|
||||
# It the latter is set to 0, the acquisition server issues its own timestamp upon reception of the stimulation.
|
||||
|
||||
import sys
|
||||
import socket
|
||||
from time import time, sleep
|
||||
|
||||
# host and port of tcp tagging server
|
||||
HOST = '127.0.0.1'
|
||||
PORT = 15361
|
||||
|
||||
# Event identifier (See stimulation codes in OpenVibe documentation)
|
||||
EVENT_ID = 5+0x8100
|
||||
|
||||
# Artificial delay (ms). It may need to be increased if the time to send the tag is too long and causes tag loss.
|
||||
DELAY=0
|
||||
|
||||
# transform a value into an array of byte values in little-endian order.
|
||||
def to_byte(value, length):
|
||||
for x in range(length):
|
||||
yield value%256
|
||||
value//=256
|
||||
|
||||
# connect
|
||||
s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
s.connect((HOST, PORT))
|
||||
|
||||
for i in range(100):
|
||||
# create the three pieces of the tag, padding, event_id and timestamp
|
||||
padding=[0]*8
|
||||
event_id=list(to_byte(EVENT_ID, 8))
|
||||
|
||||
# timestamp can be either the posix time in ms, or 0 to let the acquisition server timestamp the tag itself.
|
||||
timestamp=list(to_byte(int(time()*1000)+DELAY, 8))
|
||||
|
||||
# send tag and sleep
|
||||
s.sendall(bytearray(padding+event_id+timestamp))
|
||||
sleep(1)
|
||||
|
||||
s.close()
|
||||
|
||||
+163
@@ -0,0 +1,163 @@
|
||||
#include "ovasCPluginTCPTagging.h"
|
||||
|
||||
#include <system/ovCTime.h>
|
||||
|
||||
#include "../ovasCSettingsHelper.h"
|
||||
#include "../ovasCSettingsHelperOperators.h"
|
||||
|
||||
// #define TCPTAGGING_DEBUG
|
||||
#if defined(TCPTAGGING_DEBUG)
|
||||
#include <iomanip>
|
||||
#endif
|
||||
|
||||
namespace OpenViBE {
|
||||
namespace AcquisitionServer {
|
||||
namespace Plugins {
|
||||
|
||||
CPluginTCPTagging::CPluginTCPTagging(const Kernel::IKernelContext& ctx)
|
||||
: IAcquisitionServerPlugin(ctx, "AcquisitionServer_Plugin_TCPTagging"),
|
||||
m_port(15361)
|
||||
{
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Info << "Loading plugin: TCP Tagging\n";
|
||||
m_settings.add("TCP_Tagging_Port", &m_port);
|
||||
m_settings.load();
|
||||
}
|
||||
|
||||
bool CPluginTCPTagging::startHook(const std::vector<CString>& /*vSelectedChannelNames*/, const size_t /*sampling*/, const size_t /*nChannel*/,
|
||||
const size_t /*nSamplePerSentBlock*/)
|
||||
{
|
||||
// initialize tag stream
|
||||
// this may throw exceptions, e.g. when the port is already in use.
|
||||
try { m_scopedTagStream.reset(new CTagStream(m_port)); }
|
||||
catch (std::exception& e)
|
||||
{
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Error << "Could not create tag stream for TCP Tagging [" << e.what()
|
||||
<< "]. Make sure the port " << m_port << " is not already reserved by another Acquisition Server.\n";
|
||||
return false;
|
||||
}
|
||||
|
||||
// Initialize time counters.
|
||||
m_previousClockTime = System::Time::zgetTimeRaw(false);
|
||||
m_previousSampleTime = 0;
|
||||
m_lastTagTime = 0;
|
||||
m_lastTagTimeAdjusted = 0;
|
||||
m_warningPrinted = false;
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
void CPluginTCPTagging::stopHook() { m_scopedTagStream.reset(); }
|
||||
|
||||
// n.b. With this version of tcp tagging, all the timestamps are in fixed point
|
||||
void CPluginTCPTagging::loopHook(std::deque<std::vector<float>>& /*vPendingBuffer*/,
|
||||
CStimulationSet& stimulationSet, uint64_t /*start*/, uint64_t /*end*/, const uint64_t sampleTime)
|
||||
{
|
||||
const uint64_t clockTime = System::Time::zgetTimeRaw(false);
|
||||
|
||||
Tag tag;
|
||||
|
||||
// n.b. the last chunk received but not yet sent is between timestamps [previousClockTime,clockTime]. Note
|
||||
// that more stims may be arriving in the loop meaning that they are newer than 'clockTime'. This is not an issue,
|
||||
// they will be scheduled with future samples.
|
||||
|
||||
// Collect tags from the stream until exhaustion.
|
||||
while (m_scopedTagStream.get() && m_scopedTagStream->pop(tag))
|
||||
{
|
||||
const uint64_t tagTime = tag.timestamp;
|
||||
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Trace << "New Tag received (" << tag.flags << ", " << tag.identifier << ", "
|
||||
<< CTime(tagTime).toSeconds() << "s) at "
|
||||
<< CTime(clockTime).toSeconds() << "s\n";
|
||||
|
||||
uint64_t tagDelay = 0;
|
||||
// Check that the timestamp fits the current chunk. Unfortunately we cannot send stimulations to the past.
|
||||
if (tagTime < m_previousClockTime)
|
||||
{
|
||||
// This condition is relatively easy to achieve with high sampling rate & small block size with frequent stims (like in P300)
|
||||
tag.timestamp = m_previousClockTime;
|
||||
tagDelay = tag.timestamp - tagTime;
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Trace << "A tag "
|
||||
<< tag.identifier << " is stamped before the current chunk start; it will be late"
|
||||
<< " (delay " << CTime(tagDelay).toSeconds() * 1000.0 << "ms)\n";
|
||||
}
|
||||
if (tagTime < m_lastTagTime)
|
||||
{
|
||||
tag.timestamp = m_lastTagTime;
|
||||
tagDelay = tag.timestamp - tagTime;
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Trace << "A tag "
|
||||
<< tag.identifier << " is stamped before the previous tag; will delay"
|
||||
<< " (delay " << CTime(tagDelay).toSeconds() * 1000.0 << "ms)\n";
|
||||
}
|
||||
m_lastTagTime = tagTime;
|
||||
|
||||
// This simple and intuitive implementation has issues if the device lags:
|
||||
// if previoussampletime does not advance evenly we get problems with tag ordering; this may also be related to the frequency this function is called
|
||||
// const uint64_t tagOffsetClock = tag.timestamp - m_previousClockTime; // How far in time the marker is from the last call to this function
|
||||
// uint64_t adjustedTagTime = m_previousSampleTime + tagOffsetClock;
|
||||
|
||||
// This version estimates how far a sample is between [t1,t2] where the stamps t1,t2 are realtime of previous and current call,
|
||||
// and uses this fraction on [previousSampleTime,currentSampleTime] to get an adjustment in terms of sample time. It is
|
||||
// equivalent to interpolating between the two. To avoid implementing unsigned fixed point division, we go for doubles.
|
||||
// This should give 10e-15 decimal precision which should be more than enough for EEG.
|
||||
|
||||
const double elapsedClockTime = CTime(clockTime - m_previousClockTime).toSeconds(); // n.b. here we assume the clocks will not run backwards
|
||||
const double elapsedSampleTime = CTime(sampleTime - m_previousSampleTime).toSeconds();
|
||||
const double tagOffsetClock = CTime(tag.timestamp - m_previousClockTime).toSeconds(); // How far in time the marker is from the last call to this
|
||||
|
||||
const double scaling = elapsedSampleTime / elapsedClockTime;
|
||||
const double interpolatedOffset = (elapsedClockTime > 0 ? (tagOffsetClock * scaling) : 0);
|
||||
const uint64_t offsetSampleTime = CTime(interpolatedOffset).time();
|
||||
|
||||
uint64_t adjustedTagTime = m_previousSampleTime +
|
||||
offsetSampleTime; // Time since the beginning of the current buffer (as approx by time of the last sample of the prev. run)
|
||||
|
||||
if (adjustedTagTime < m_lastTagTimeAdjusted)
|
||||
{
|
||||
tagDelay = m_lastTagTimeAdjusted - adjustedTagTime;
|
||||
m_kernelCtx.getLogManager() << Kernel::LogLevel_Trace << "A tag "
|
||||
<< tag.identifier << " was adjusted before the previous tag; will delay"
|
||||
<< " (delay " << CTime(tagDelay).toSeconds() * 1000.0 << "ms)"
|
||||
<< " oc " << tagOffsetClock * 1000.0 << "ms\n";
|
||||
adjustedTagTime = m_lastTagTimeAdjusted;
|
||||
}
|
||||
m_lastTagTimeAdjusted = adjustedTagTime;
|
||||
|
||||
#if defined(TCPTAGGING_DEBUG)
|
||||
// If the amp and the AS computer are both behaving similarly, the scaling term should be very
|
||||
// close to 1 and the diff term should be nearly zero. In that case the approach behaves
|
||||
// like the simple, commented out solution above.
|
||||
std::cout << "Set tag " << tag.identifier
|
||||
<< " at " << std::setprecision(6) << CTime(adjustedTagTime).toSeconds()
|
||||
<< " (pst = " << CTime(m_previousSampleTime).toSeconds() << "s,"
|
||||
<< " pct = " << CTime(m_previousClockTime).toSeconds() << "s,"
|
||||
<< " otag = " << CTime(tagTime).toSeconds() << "s,"
|
||||
<< " ntag = " << CTime(tag.timestamp).toSeconds() << "s,"
|
||||
<< " s = " << scaling << ","
|
||||
<< " off = " << tagOffsetClock << "s,"
|
||||
<< " offI = " << interpolatedOffset << "s,"
|
||||
<< " diff = " << (interpolatedOffset - tagOffsetClock)*1000.0 << "ms,"
|
||||
<< " del = " << CTime(tagDelay).toSeconds()*1000.0 << "ms"
|
||||
<< ")\n";
|
||||
#endif
|
||||
|
||||
if (tagDelay > 0)
|
||||
{
|
||||
// Indicates that the next tag after this one may not be correctly placed in time.
|
||||
// The duration encodes our estimate how much the tag was delayed. This has
|
||||
// the benefit that this knowledge can be inserted into file recordings and is not lost like logs potentially.
|
||||
stimulationSet.appendStimulation(OVTK_GDF_Incorrect, adjustedTagTime, tagDelay);
|
||||
}
|
||||
|
||||
// Insert tag into the stimulation set.
|
||||
stimulationSet.appendStimulation(tag.identifier, adjustedTagTime, 0);
|
||||
}
|
||||
|
||||
// Update time counters. Basically these counters allow to map the time a stamp was received to the time related to the sample buffers,
|
||||
// as we know this function is called right after receiving samples from a device.
|
||||
m_previousClockTime = clockTime;
|
||||
m_previousSampleTime = sampleTime;
|
||||
}
|
||||
|
||||
} // namespace Plugins
|
||||
} // namespace AcquisitionServer
|
||||
} // namespace OpenViBE
|
||||
+57
@@ -0,0 +1,57 @@
|
||||
#pragma once
|
||||
|
||||
/**
|
||||
* \brief Acquisition Server plugin adding the capability to receive stimulations from external sources
|
||||
* via TCP/IP.
|
||||
*
|
||||
* The stimulation format is the same as with Shared Memory Tagging. It comprises three blocks of 8 bytes:
|
||||
*
|
||||
* ----------------------------------------------------------------------
|
||||
* | padding (8 bytes) | event id (8 bytes) | timestamp (8 bytes) |
|
||||
* ----------------------------------------------------------------------
|
||||
*
|
||||
* The padding is only for consistency with Shared Memory Tagging and has no utility.
|
||||
* The event id informs about the type of event happening.
|
||||
* The timestamp is the posix time (ms since Epoch) at the moment of the event.
|
||||
* It the latter is set to 0, the acquisition server issues its own timestamp upon reception of the stimulation.
|
||||
*
|
||||
* Have a look at contrib/plugins/server-extensions/tcp-tagging/client-example to learn about the protocol
|
||||
* to send stimulations from the client.
|
||||
*/
|
||||
|
||||
#include "ovasIAcquisitionServerPlugin.h"
|
||||
#include "ovasCTagStream.h"
|
||||
|
||||
namespace OpenViBE {
|
||||
namespace AcquisitionServer {
|
||||
namespace Plugins {
|
||||
class CPluginTCPTagging final : public IAcquisitionServerPlugin
|
||||
{
|
||||
public:
|
||||
explicit CPluginTCPTagging(const Kernel::IKernelContext& ctx);
|
||||
~CPluginTCPTagging() override { }
|
||||
|
||||
// Overrides virtual method startHook inherited from class IAcquisitionServerPlugin.
|
||||
bool startHook(const std::vector<CString>& vSelectedChannelNames, const size_t sampling, const size_t nChannel, const size_t nSamplePerSentBlock) override;
|
||||
|
||||
// Overrides virtual method stopHook inherited from class IAcquisitionServerPlugin
|
||||
void stopHook() override;
|
||||
|
||||
// Overrides virtual method loopHook inherited from class IAcquisitionServerPlugin.
|
||||
void loopHook(std::deque<std::vector<float>>& vPendingBuffer, CStimulationSet& stimulationSet, const uint64_t start, const uint64_t end,
|
||||
const uint64_t sampleTime) override;
|
||||
|
||||
private:
|
||||
uint64_t m_previousClockTime = 0;
|
||||
uint64_t m_previousSampleTime = 0;
|
||||
uint64_t m_lastTagTime = 0;
|
||||
uint64_t m_lastTagTimeAdjusted = 0;
|
||||
|
||||
std::unique_ptr<CTagStream> m_scopedTagStream;
|
||||
size_t m_port = 0;
|
||||
|
||||
bool m_warningPrinted = false;
|
||||
};
|
||||
} // namespace Plugins
|
||||
} // namespace AcquisitionServer
|
||||
} // namespace OpenViBE
|
||||
Masterarbeit/openvibe/extras-master/contrib/plugins/server-extensions/tcp-tagging/ovasCTagStream.cpp
Executable
+131
@@ -0,0 +1,131 @@
|
||||
#include "ovasCTagStream.h"
|
||||
|
||||
#include <system/ovCTime.h>
|
||||
#include <boost/bind.hpp>
|
||||
|
||||
#include <iostream>
|
||||
|
||||
#include <thread>
|
||||
#include <mutex>
|
||||
|
||||
namespace OpenViBE {
|
||||
namespace AcquisitionServer {
|
||||
namespace Plugins {
|
||||
|
||||
void CTagQueue::push(const Tag& tag)
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(m_mutex);
|
||||
m_queue.push(tag);
|
||||
}
|
||||
|
||||
bool CTagQueue::pop(Tag& tag)
|
||||
{
|
||||
std::lock_guard<std::mutex> guard(m_mutex);
|
||||
if (m_queue.empty()) { return false; }
|
||||
tag = m_queue.front();
|
||||
m_queue.pop();
|
||||
return true;
|
||||
}
|
||||
|
||||
void CTagSession::start()
|
||||
{
|
||||
m_errorState = 0;
|
||||
startRead();
|
||||
}
|
||||
|
||||
void CTagSession::startRead()
|
||||
{
|
||||
// Caveat: a shared pointer is used (instead of simply using this) to ensure that this instance of TagSession is still alive when the callback is called.
|
||||
async_read(m_socket, boost::asio::buffer(static_cast<void*>(&m_tag), sizeof(Tag)), boost::bind(&CTagSession::handleRead, shared_from_this(), _1));
|
||||
}
|
||||
|
||||
void CTagSession::handleRead(const boost::system::error_code& error)
|
||||
{
|
||||
if (!error)
|
||||
{
|
||||
if (m_tag.timestamp == 0 || (m_tag.flags & FLAG_AUTOSTAMP_SERVERSIDE))
|
||||
{
|
||||
// Client didn't provide timestamp or asked the server to do it. Stamp current time.
|
||||
m_tag.timestamp = System::Time::zgetTimeRaw(false);
|
||||
}
|
||||
else if (!(m_tag.flags & FLAG_FPTIME))
|
||||
{
|
||||
// Client provided stamp but not in FPTIME
|
||||
m_tag.timestamp = System::Time::zgetTimeRaw(false);
|
||||
if (!(m_errorState & (1LL << 1)))
|
||||
{
|
||||
// @fixme not appropriate to print errors from a thread, but better than silent fail
|
||||
std::cout << "[WARNING] TCP Tagging: Received tag(s) not in fixed point time. Not supported, will replace with server time.\n";
|
||||
m_errorState |= (1LL << 1);
|
||||
}
|
||||
}
|
||||
|
||||
// Push tag to the queue.
|
||||
m_queuePtr->push(m_tag);
|
||||
|
||||
// Continue reading.
|
||||
startRead();
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
CTagServer::CTagServer(const SharedQueuePtr& queue, int port)
|
||||
: m_acceptor(m_ioService), m_queuePtr(queue)
|
||||
{
|
||||
boost::asio::ip::tcp::endpoint endp = boost::asio::ip::tcp::endpoint(boost::asio::ip::tcp::v4(), port);
|
||||
m_acceptor.open(endp.protocol());
|
||||
|
||||
// Try to make sure that the port cannot be used by multiple processes
|
||||
boost::asio::socket_base::reuse_address option(false);
|
||||
m_acceptor.set_option(option);
|
||||
|
||||
m_acceptor.bind(endp);
|
||||
m_acceptor.listen();
|
||||
}
|
||||
|
||||
|
||||
void CTagServer::run()
|
||||
{
|
||||
try
|
||||
{
|
||||
startAccept();
|
||||
m_ioService.run();
|
||||
}
|
||||
catch (std::exception&)
|
||||
{
|
||||
// TODO: log error message (needs to be thread-safe)
|
||||
}
|
||||
}
|
||||
|
||||
void CTagServer::startAccept()
|
||||
{
|
||||
SharedSessionPtr newSession(new CTagSession(m_ioService, m_queuePtr));
|
||||
// Note: if this instance of CTagSever is destroyed then the associated io_service is destroyed as well.
|
||||
// Therefore the call-back will never be called if this instance is destroyed and it is safe to use this instead of a shared pointer.
|
||||
|
||||
m_acceptor.async_accept(newSession->socket(), boost::bind(&CTagServer::handleAccept, this, newSession, _1));
|
||||
}
|
||||
|
||||
void CTagServer::handleAccept(SharedSessionPtr& session, const boost::system::error_code& error)
|
||||
{
|
||||
if (!error) { session->start(); }
|
||||
startAccept();
|
||||
}
|
||||
|
||||
CTagStream::CTagStream(const int port) : m_queuePtr(new CTagQueue), m_port(port)
|
||||
{
|
||||
// can throw exceptions, e.g. when the port is already in use.
|
||||
m_serverPtr.reset(new CTagServer(m_queuePtr, m_port));
|
||||
m_threadPtr.reset(new std::thread(&CTagStream::startServer, this));
|
||||
}
|
||||
|
||||
CTagStream::~CTagStream()
|
||||
{
|
||||
// m_serverPtr and m_threadPtr cannot be null
|
||||
m_serverPtr->stop();
|
||||
m_threadPtr->join();
|
||||
}
|
||||
|
||||
} // namespace Plugins
|
||||
} // namespace AcquisitionServer
|
||||
} // namespace OpenViBE
|
||||
Executable
+121
@@ -0,0 +1,121 @@
|
||||
#pragma once
|
||||
|
||||
#include <queue>
|
||||
#include <boost/asio.hpp>
|
||||
|
||||
#include <mutex>
|
||||
#include <thread>
|
||||
|
||||
// PluginTCPTagging relies on four auxilliary classes: CTagQueue, CTagSession, CTagServer and CTagStream.
|
||||
// CTagQueue implements a trivial queue to store tags with exclusive locking.
|
||||
// CTagServer implements a server that simply binds to a port and waits for incoming connections.
|
||||
// CTagSession represents an individual connection with a client and holds a connection handle (socket)
|
||||
// and a data buffer to store incoming data.
|
||||
// The use of shared pointers is instrumental to ensure that instances are still alive when call-backs are
|
||||
// called and avoid memory corruption.
|
||||
// The CTagStream class implements a stream to allow to collect tags. Upon instantiation, it creates an instance
|
||||
// of CTagServer and starts the server in an auxilliary thread.
|
||||
// The exchange of data between the main tread and the auxilliary thread is performed via a lockfree queue (boost).
|
||||
|
||||
namespace OpenViBE {
|
||||
namespace AcquisitionServer {
|
||||
namespace Plugins {
|
||||
|
||||
// A Tag consists of an identifier to inform about the type of event
|
||||
// and a timestamp corresponding to the time at which the event occurrs.
|
||||
struct Tag
|
||||
{
|
||||
Tag(): flags(0), identifier(0), timestamp(0) {}
|
||||
uint64_t flags, identifier, timestamp;
|
||||
};
|
||||
|
||||
// Note: duplicated in TCP Tagging module in openvibe
|
||||
enum TCP_Tagging_Flags
|
||||
{
|
||||
FLAG_FPTIME = (1LL << 0), // The time given is fixed point time.
|
||||
FLAG_AUTOSTAMP_CLIENTSIDE = (1LL << 1), // Ignore given stamp, bake timestamp on client side before sending
|
||||
FLAG_AUTOSTAMP_SERVERSIDE = (1LL << 2) // Ignore given stamp, bake timestamp on server side when receiving
|
||||
};
|
||||
|
||||
class CTagSession; // forward declaration of CTagSession to define SharedSessionPtr
|
||||
class CTagQueue; // forward declaration of CTagQueue to define SharedQueuePtr
|
||||
class CTagServer; // forward declaration of CTagServer to define ScopedServerPtr
|
||||
|
||||
typedef std::shared_ptr<CTagQueue> SharedQueuePtr;
|
||||
typedef std::shared_ptr<CTagSession> SharedSessionPtr;
|
||||
typedef std::unique_ptr<CTagServer> ScopedServerPtr;
|
||||
typedef std::unique_ptr<std::thread> ScopedThreadPtr;
|
||||
|
||||
// A trivial implementation of a queue to store Tags with exclusive locking
|
||||
class CTagQueue
|
||||
{
|
||||
public:
|
||||
CTagQueue() { }
|
||||
void push(const Tag& tag);
|
||||
bool pop(Tag& tag);
|
||||
private:
|
||||
std::queue<Tag> m_queue;
|
||||
std::mutex m_mutex;
|
||||
};
|
||||
|
||||
// An instance of CTagSession is associated to every client connecting to the Tagging Server.
|
||||
// It contains a connection handle and data buffer.
|
||||
class CTagSession : public std::enable_shared_from_this<CTagSession>
|
||||
{
|
||||
public:
|
||||
CTagSession(boost::asio::io_service& ioService, const SharedQueuePtr& queue) : m_socket(ioService), m_queuePtr(queue) { }
|
||||
|
||||
boost::asio::ip::tcp::socket& socket() { return m_socket; }
|
||||
void start();
|
||||
void startRead();
|
||||
void handleRead(const boost::system::error_code& error);
|
||||
|
||||
private:
|
||||
Tag m_tag;
|
||||
boost::asio::ip::tcp::socket m_socket;
|
||||
SharedQueuePtr m_queuePtr;
|
||||
uint64_t m_errorState = 0;
|
||||
};
|
||||
|
||||
// CTagServer implements a server that binds to a port and accepts new connections.
|
||||
class CTagServer
|
||||
{
|
||||
public:
|
||||
explicit CTagServer(const SharedQueuePtr& queue, int port = 15361);
|
||||
~CTagServer() { }
|
||||
|
||||
void run();
|
||||
void stop() { m_ioService.stop(); }
|
||||
|
||||
private:
|
||||
void startAccept();
|
||||
void handleAccept(SharedSessionPtr& session, const boost::system::error_code& error);
|
||||
|
||||
boost::asio::io_service m_ioService;
|
||||
boost::asio::ip::tcp::acceptor m_acceptor;
|
||||
const SharedQueuePtr& m_queuePtr;
|
||||
};
|
||||
|
||||
// CTagStream allows to collect tags received via TCP.
|
||||
class CTagStream
|
||||
{
|
||||
// Initial memory allocation of lockfree queue.
|
||||
enum { ALLOCATE = 128 };
|
||||
|
||||
public:
|
||||
explicit CTagStream(int port = 15361);
|
||||
~CTagStream();
|
||||
|
||||
bool pop(Tag& tag) { return m_queuePtr->pop(tag); }
|
||||
|
||||
private:
|
||||
void startServer() { m_serverPtr->run(); }
|
||||
|
||||
SharedQueuePtr m_queuePtr;
|
||||
ScopedServerPtr m_serverPtr;
|
||||
ScopedThreadPtr m_threadPtr;
|
||||
int m_port = 0;
|
||||
};
|
||||
} // namespace Plugins
|
||||
} // namespace AcquisitionServer
|
||||
} // namespace OpenViBE
|
||||
+27
@@ -0,0 +1,27 @@
|
||||
PROJECT(test_tagstream)
|
||||
|
||||
IF(WIN32)
|
||||
ADD_DEFINITIONS(-DTARGET_OS_Windows)
|
||||
ENDIF(WIN32)
|
||||
IF(UNIX)
|
||||
ADD_DEFINITIONS(-DTARGET_OS_Linux)
|
||||
ENDIF(UNIX)
|
||||
|
||||
INCLUDE_DIRECTORIES(../)
|
||||
ADD_EXECUTABLE(${PROJECT_NAME} test_tagstream.cpp ../ovasCTagStream.cpp)
|
||||
SET_PROPERTY(TARGET ${PROJECT_NAME} PROPERTY FOLDER ${TESTS_FOLDER}) # Place project in folder unit-test (for some IDE)
|
||||
|
||||
INCLUDE("FindOpenViBE")
|
||||
INCLUDE("FindOpenViBEModuleSystem") # Time getter from here in the future
|
||||
INCLUDE("FindThirdPartyBoost")
|
||||
INCLUDE("FindThirdPartyBoost_System")
|
||||
INCLUDE("FindThirdPartyBoost_Thread")
|
||||
|
||||
# Unfortunately we need to install the tests as any application to find .dll/.so files
|
||||
# on both Windows and Linux.
|
||||
OV_INSTALL_LAUNCH_SCRIPT(SCRIPT_PREFIX "${PROJECT_NAME}" EXECUTABLE_NAME "${PROJECT_NAME}")
|
||||
INSTALL(TARGETS ${PROJECT_NAME}
|
||||
RUNTIME DESTINATION ${DIST_BINDIR}
|
||||
LIBRARY DESTINATION ${DIST_LIBDIR}
|
||||
ARCHIVE DESTINATION ${DIST_LIBDIR})
|
||||
|
||||
+18
@@ -0,0 +1,18 @@
|
||||
# Basic Template Test for automatic run a scenario that produce a file to be compared to a reference file
|
||||
# You need to set the name of the test according to name of scenario file and reference file
|
||||
|
||||
# Test TagStream
|
||||
|
||||
SET(TEST_NAME "TagStream")
|
||||
|
||||
IF(WIN32)
|
||||
SET(EXT cmd)
|
||||
SET(OS_FLAGS "--no-pause")
|
||||
ELSE(WIN32)
|
||||
SET(EXT sh)
|
||||
SET(OS_FLAGS "")
|
||||
ENDIF(WIN32)
|
||||
|
||||
ADD_TEST(run_${TEST_NAME} "$ENV{OV_BINARY_PATH}/test_tagstream.${EXT}" ${OS_FLAGS})
|
||||
|
||||
|
||||
+26
@@ -0,0 +1,26 @@
|
||||
#include "../ovasCTagStream.h"
|
||||
|
||||
#include <iostream>
|
||||
|
||||
int main()
|
||||
{
|
||||
bool ok = false;
|
||||
|
||||
OpenViBE::AcquisitionServer::Plugins::CTagStream tagStream1;
|
||||
|
||||
// The construction of the second TagStream must fail because of port already in use.
|
||||
try { OpenViBE::AcquisitionServer::Plugins::CTagStream tagStream2; }
|
||||
catch (std::exception&) { ok = true; } // This exception is expected, don't print
|
||||
|
||||
// The construction must succeed because another port is used.
|
||||
try { OpenViBE::AcquisitionServer::Plugins::CTagStream tagStream3(15362); }
|
||||
catch (std::exception& e)
|
||||
{
|
||||
// Unexpected exception
|
||||
std::cout << "Exception: " << e.what() << "\n";
|
||||
ok = false;
|
||||
}
|
||||
|
||||
if (!ok) { return 1; }
|
||||
return 0;
|
||||
}
|
||||
Reference in New Issue
Block a user