ddt 1.4.0
 
Loading...
Searching...
No Matches
ddtMemoryAccessor.hpp
Go to the documentation of this file.
1
10
11#ifndef DDTMEMORYACCESSOR_H_
12#define DDTMEMORYACCESSOR_H_
13
14#include <boost/circular_buffer.hpp>
15#include <boost/interprocess/containers/string.hpp>
16#include <boost/interprocess/containers/vector.hpp>
17#include <boost/interprocess/managed_shared_memory.hpp>
18#include <boost/signals2/signal.hpp>
19#include <fstream>
20#include <future>
21#include <iostream>
22
23#include "ddt/ddtConstants.hpp"
24#include "ddt/ddtCrc32.hpp"
25#include "ddt/ddtLogger.hpp"
26
27namespace ip = boost::interprocess;
28
32typedef ip::managed_shared_memory::segment_manager segment_manager_t;
33
37typedef ip::allocator<void, segment_manager_t> void_allocator;
38
42typedef ip::allocator<uint8_t, segment_manager_t> uint8_allocator;
43
47typedef ip::vector<uint8_t, uint8_allocator> uint8_vector;
48
52typedef ip::allocator<uint16_t, segment_manager_t> uint16_allocator;
53
57typedef ip::vector<uint16_t, uint16_allocator> uint16_vector;
58
62typedef ip::allocator<char, segment_manager_t> char_allocator;
63
67typedef ip::basic_string<char, std::char_traits<char>, char_allocator>
69
73typedef boost::signals2::signal<void()> SignalT;
74
75namespace ddt {
76
84 int32_t topic_id = 0;
85
90
95
99 int32_t sample_id;
100
105
113 DataSampleShared(const int32_t id, const int md_length,
114 const int vector_length, const void_allocator &void_alloc)
115 : meta_data_length(md_length),
116 meta_data(md_length, uint8_t(), void_alloc),
117 sample_id(id),
118 data(vector_length, uint8_t(), void_alloc) {}
119};
120
131
135 uint32_t checksum;
136
141
145 int64_t writer_index = -1;
146
150 uint64_t timestamp = 0;
151
156
164 DataPacketShared(const char *const ds_id, const int32_t check,
165 const int vector_length, const void_allocator &void_alloc)
166 : data_stream_identifier(ds_id, void_alloc),
167 checksum(check),
168 sample_length(vector_length),
169 sample(0, META_DATA_LENGTH, vector_length, void_alloc) {}
170};
171
179 int32_t topic_id = 0;
180
185
189 std::vector<uint8_t> meta_data;
190
194 int32_t sample_id;
195
199 std::vector<uint8_t> data;
200
207 DataSample(const int32_t id, const int md_length, const int vector_length)
208 : meta_data_length(md_length),
209 meta_data(md_length),
210 sample_id(id),
211 data(vector_length) {}
212};
213
222
226 uint32_t checksum;
227
232
236 int64_t writer_index = -1;
237
241 uint64_t timestamp = 0;
242
247
254 DataPacket(const char *const ds_id, const int32_t check,
255 const int vector_length)
256 : data_stream_identifier(ds_id),
257 checksum(check),
258 sample_length(vector_length),
259 sample(0, META_DATA_LENGTH, vector_length) {}
260};
261
266 public:
271
280 explicit DdtMemoryAccessor(const std::string &shm_id,
281 const std::string &data_stream_identifier,
282 DdtLogger *logger, const uint64_t time_window = 0,
283 const int32_t reading_interval = 10);
284
288 virtual ~DdtMemoryAccessor();
289
295 const uint32_t ComputeChecksum(DataSampleShared *const data_sample_shared);
296
302 const uint32_t ComputeChecksum(DataSample *const data_sample);
303
308 int32_t OpenSharedMemory();
309
313 void CloseSharedMemory();
314
326 void WriteData(const int32_t writer_index, const int32_t topic_id,
327 const int32_t sample_id, const uint8_t *datavec,
328 const int32_t datavec_size, const uint8_t *metadata_vec,
329 const int32_t metadatavec_size, const uint64_t timestamp);
330
334 void StartReading();
335
339 void StopReading();
340
346 void SetSizeConstraints(const int32_t max_sample_size, const int32_t space);
347
351 void Reattach();
352
358 void NewData();
359
369 void get_data_packet(std::string *stream_identifier, uint32_t *checksum,
370 int32_t *sample_length, int64_t *writer_idx,
371 uint64_t *timestamp, DataSample **sample);
372
377 bool get_data_available();
378
384
390
395 void Reset();
396
401 void set_pub_unreg(const bool STATE);
402
408 void set_compute_checksum(const bool compute_crc);
409
414 bool get_compute_checksum() const;
415
420 bool get_is_initialized() const;
421
422 private:
423 SignalT data_available_signal;
424
428 void ReadData();
429
433 void SetDefaults();
434
443 void Init(const std::string &mem_id, const std::string &stream_id,
444 DdtLogger *ddt_logger, const uint64_t time_win,
445 const int32_t interval);
446
450 void PrintData();
451
456 int32_t CreateNewShm();
457
462 int32_t SearchCircBuffer();
463
468 int32_t SearchWriterIndex();
469
470 ip::managed_shared_memory *managed_shm;
471
476 typedef ip::allocator<DataPacketShared,
477 ip::managed_shared_memory::segment_manager>
478 cb_alloc;
479
483 typedef boost::circular_buffer<DataPacketShared, cb_alloc> cb;
484
485 std::string shm_id;
486 std::string data_stream_identifier;
487 uint64_t time_window; // in [ms]
488 int32_t reading_interval; // in [ms]
489 cb *circ_buffer;
490 std::atomic<int64_t> *writer_index;
491 int32_t local_index;
492 int64_t reader_index;
493
494 int32_t number_of_unread_elements;
495 int32_t circ_buf_capacity;
496 int32_t number_of_lost_packages;
497
498 std::mutex circ_buffer_mutex;
499 std::mutex packets_mutex;
500
501 bool is_initialized;
502
503 std::promise<void> exit_signal;
504 std::future<void> future_object;
505
506 std::atomic<bool> reading_active;
507 std::atomic<bool> pub_unreg;
508 std::atomic<bool> compute_checksum;
509
510 int32_t max_data_sample_size;
511 int additional_space;
512
513 std::list<DataPacketShared *> packets;
514
515 DdtLogger *logger;
516};
517
518} // namespace ddt
519
520#endif /* DDTMEMORYACCESSOR_H_ */
Definition ddtLogger.hpp:43
virtual ~DdtMemoryAccessor()
Definition ddtMemoryAccessor.cpp:28
void set_compute_checksum(const bool compute_crc)
Definition ddtMemoryAccessor.cpp:542
void StopReading()
Definition ddtMemoryAccessor.cpp:317
bool get_is_initialized() const
Definition ddtMemoryAccessor.cpp:559
bool get_data_available()
Definition ddtMemoryAccessor.cpp:426
void set_pub_unreg(const bool STATE)
Definition ddtMemoryAccessor.cpp:540
void CloseSharedMemory()
Definition ddtMemoryAccessor.cpp:73
void StartReading()
Definition ddtMemoryAccessor.cpp:302
int32_t OpenSharedMemory()
Definition ddtMemoryAccessor.cpp:117
void NewData()
Definition ddtMemoryAccessor.cpp:345
int32_t get_number_of_unread_elements()
Definition ddtMemoryAccessor.cpp:436
void WriteData(const int32_t writer_index, const int32_t topic_id, const int32_t sample_id, const uint8_t *datavec, const int32_t datavec_size, const uint8_t *metadata_vec, const int32_t metadatavec_size, const uint64_t timestamp)
Definition ddtMemoryAccessor.cpp:183
void Reattach()
Definition ddtMemoryAccessor.cpp:337
void Reset()
Definition ddtMemoryAccessor.cpp:527
void get_data_packet(std::string *stream_identifier, uint32_t *checksum, int32_t *sample_length, int64_t *writer_idx, uint64_t *timestamp, DataSample **sample)
Definition ddtMemoryAccessor.cpp:441
DdtMemoryAccessor()
Definition ddtMemoryAccessor.cpp:17
SignalT * DataAvailableSignal()
Definition ddtMemoryAccessor.cpp:523
bool get_compute_checksum() const
Definition ddtMemoryAccessor.cpp:555
void SetSizeConstraints(const int32_t max_sample_size, const int32_t space)
Definition ddtMemoryAccessor.cpp:331
const uint32_t ComputeChecksum(DataSampleShared *const data_sample_shared)
Definition ddtMemoryAccessor.cpp:80
Contains common used constants. This file shall contain constants that can be used by all application...
Class to wrap the usage of log4cplus as logging utility. This file provides a wrapper class for the u...
ip::allocator< char, segment_manager_t > char_allocator
Definition ddtMemoryAccessor.hpp:62
ip::vector< uint8_t, uint8_allocator > uint8_vector
Definition ddtMemoryAccessor.hpp:47
ip::allocator< uint16_t, segment_manager_t > uint16_allocator
Definition ddtMemoryAccessor.hpp:52
ip::allocator< uint8_t, segment_manager_t > uint8_allocator
Definition ddtMemoryAccessor.hpp:42
ip::vector< uint16_t, uint16_allocator > uint16_vector
Definition ddtMemoryAccessor.hpp:57
ip::allocator< void, segment_manager_t > void_allocator
Definition ddtMemoryAccessor.hpp:37
ip::basic_string< char, std::char_traits< char >, char_allocator > char_string
Definition ddtMemoryAccessor.hpp:68
boost::signals2::signal< void()> SignalT
Definition ddtMemoryAccessor.hpp:73
ip::managed_shared_memory::segment_manager segment_manager_t
Definition ddtMemoryAccessor.hpp:32
Definition ddtClient.hpp:32
const int META_DATA_LENGTH
Definition ddtConstants.hpp:59
Definition ddtMemoryAccessor.hpp:126
DataPacketShared(const char *const ds_id, const int32_t check, const int vector_length, const void_allocator &void_alloc)
Definition ddtMemoryAccessor.hpp:164
int32_t sample_length
Definition ddtMemoryAccessor.hpp:140
DataSampleShared sample
Definition ddtMemoryAccessor.hpp:155
uint32_t checksum
Definition ddtMemoryAccessor.hpp:135
uint64_t timestamp
Definition ddtMemoryAccessor.hpp:150
int64_t writer_index
Definition ddtMemoryAccessor.hpp:145
char_string data_stream_identifier
Definition ddtMemoryAccessor.hpp:130
int32_t sample_length
Definition ddtMemoryAccessor.hpp:231
DataPacket(const char *const ds_id, const int32_t check, const int vector_length)
Definition ddtMemoryAccessor.hpp:254
uint32_t checksum
Definition ddtMemoryAccessor.hpp:226
DataSample sample
Definition ddtMemoryAccessor.hpp:246
std::string data_stream_identifier
Definition ddtMemoryAccessor.hpp:221
int64_t writer_index
Definition ddtMemoryAccessor.hpp:236
uint64_t timestamp
Definition ddtMemoryAccessor.hpp:241
Definition ddtMemoryAccessor.hpp:80
int32_t topic_id
Definition ddtMemoryAccessor.hpp:84
int32_t sample_id
Definition ddtMemoryAccessor.hpp:99
int32_t meta_data_length
Definition ddtMemoryAccessor.hpp:89
uint8_vector meta_data
Definition ddtMemoryAccessor.hpp:94
uint8_vector data
Definition ddtMemoryAccessor.hpp:104
DataSampleShared(const int32_t id, const int md_length, const int vector_length, const void_allocator &void_alloc)
Definition ddtMemoryAccessor.hpp:113
Definition ddtMemoryAccessor.hpp:175
std::vector< uint8_t > meta_data
Definition ddtMemoryAccessor.hpp:189
int32_t topic_id
Definition ddtMemoryAccessor.hpp:179
DataSample(const int32_t id, const int md_length, const int vector_length)
Definition ddtMemoryAccessor.hpp:207
std::vector< uint8_t > data
Definition ddtMemoryAccessor.hpp:199
int32_t sample_id
Definition ddtMemoryAccessor.hpp:194
int32_t meta_data_length
Definition ddtMemoryAccessor.hpp:184