12#ifndef RTCTK_COMPONENTFRAMEWORK_TEST_REPOSITORYSUBSCRIBERIFTESTSUITE_HPP
13#define RTCTK_COMPONENTFRAMEWORK_TEST_REPOSITORYSUBSCRIBERIFTESTSUITE_HPP
18#include <gtest/gtest.h>
26static std::shared_ptr<RepositorySubscriberIf> MakeRepository();
32 std::scoped_lock lock{m_mutex};
33 m_data.push_back(path);
37 std::scoped_lock lock{m_mutex};
38 return m_data.at(idx);
42 std::scoped_lock lock{m_mutex};
47 std::scoped_lock lock{m_mutex};
48 auto it = std::find(m_data.begin(), m_data.end(), val);
49 return it != m_data.end();
52 void AwaitSize(
size_t target_size, std::chrono::seconds timeout = std::chrono::seconds(5)) {
53 auto t_start = std::chrono::steady_clock::now();
54 while (
Size() < target_size) {
55 std::this_thread::sleep_for(std::chrono::microseconds(100));
56 if ((std::chrono::steady_clock::now() - t_start) > timeout) {
57 FAIL() <<
"ThreadSafeQ::AwaitSize: Timed out waiting for size "
58 << std::to_string(target_size);
65 std::vector<T> m_data;
79 repo = MakeRepository();
80 path1 =
"/foo"_dppath;
81 path2 =
"/bar"_dppath;
82 link =
"/link"_dppath;
86 if (
repo ==
nullptr) {
102 std::shared_ptr<RepositorySubscriberIf>
repo;
114 repo->SendRequest(req).Wait();
124 auto sub1 = repo->Subscribe<
int>(
131 [&](std::exception_ptr error) { error_q_1.
PushBack(error); });
133 auto sub2 = repo->Subscribe(
148 repo->SendRequest(req).Wait();
153 EXPECT_TRUE(path_q_1.
Contains(path1));
154 EXPECT_TRUE(value_q_1.
Contains(int_write));
155 EXPECT_TRUE(path_q_2.
Contains(path2));
156 EXPECT_EQ(error_q_1.
Size(), 0);
162 num_cb_q_1 = path_q_1.
Size();
163 num_cb_q_2 = path_q_2.
Size();
170 repo->SendRequest(req).Wait();
173 std::this_thread::sleep_for(std::chrono::milliseconds(200));
176 EXPECT_EQ(error_q_1.
Size(), 0);
177 EXPECT_EQ(path_q_1.
Size(), num_cb_q_1);
178 EXPECT_EQ(value_q_1.
Size(), num_cb_q_1);
179 EXPECT_EQ(path_q_2.
Size(), num_cb_q_2);
184 repo->CreateDataPoint(path1, 0);
191 auto sub1 = repo->Subscribe(
199 auto sub2 = repo->Subscribe(
210 repo->WriteDataPoint(path1, 3);
214 EXPECT_TRUE(path_q_1.
Contains(path1));
215 EXPECT_TRUE(path_q_2.
Contains(path1));
219 num_cb_q_1 = path_q_1.
Size();
220 num_cb_q_2 = path_q_2.
Size();
223 repo->WriteDataPoint(path1, 4);
228 std::this_thread::sleep_for(std::chrono::milliseconds(200));
230 EXPECT_EQ(path_q_1.
Size(), (num_cb_q_1 + 1));
231 EXPECT_EQ(path_q_2.
Size(), num_cb_q_2);
240 repo->SendRequest(req).Wait();
243 std::atomic_bool stop =
false;
246 auto writer_func = [&] {
248 while (stop ==
false) {
249 repo->WriteDataPoint(path1, value);
255 auto subscriber_func = [&] {
256 for (
unsigned i = 0; i < 1000; i++) {
259 auto sub = repo->Subscribe(
273 std::thread t1{writer_func};
274 std::thread t2{subscriber_func};
285 repo->SendRequest(req).Wait();
293 auto sub = repo->Subscribe<
int>(
302 [&](std::exception_ptr error) { error_q.
PushBack(error); });
310 repo->WriteDataPoint(path1, int_write);
315 EXPECT_TRUE(value_q.
Contains(int_write));
316 EXPECT_EQ(error_q.
Size(), 0);
322 repo->WriteDataPoint(link, int_write);
327 EXPECT_TRUE(value_q.
Contains(int_write));
328 EXPECT_EQ(error_q.
Size(), 0);
335 int num_cb_q = path_q.
Size();
338 repo->WriteDataPoint(path1, int_write);
341 std::this_thread::sleep_for(std::chrono::milliseconds(200));
343 EXPECT_EQ(error_q.
Size(), 0);
344 EXPECT_EQ(path_q.
Size(), num_cb_q);
345 EXPECT_EQ(value_q.
Size(), num_cb_q);
355 repo->SendRequest(req).Wait();
362 auto sub = repo->Subscribe<
int>(
369 [&](std::exception_ptr error) { error_q.
PushBack(error); });
376 repo->WriteDataPoint(path1, int_write);
381 EXPECT_TRUE(value_q.
Contains(int_write));
382 EXPECT_EQ(error_q.
Size(), 0);
385 int num_cb_q = path_q.
Size();
386 repo->DeleteDataPoint(path1);
389 std::this_thread::sleep_for(std::chrono::milliseconds(200));
391 EXPECT_EQ(error_q.
Size(), 0);
392 EXPECT_EQ(path_q.
Size(), num_cb_q);
393 EXPECT_EQ(value_q.
Size(), num_cb_q);
399 int int_create_2 = 100;
404 repo->SendRequest(req).Wait();
411 auto sub = repo->Subscribe<
int>(
418 [&](std::exception_ptr error) { error_q.
PushBack(error); });
425 repo->WriteDataPoint(path1, int_write);
430 EXPECT_TRUE(value_q.
Contains(int_write));
431 EXPECT_EQ(error_q.
Size(), 0);
437 repo->SendRequest(req).Wait();
444 EXPECT_EQ(error_q.
Size(), 0);
447 int num_cb_q = path_q.
Size();
451 repo->SendRequest(req).Wait();
454 repo->WriteDataPoint(path1, int_write);
455 repo->WriteDataPoint(path2, int_write);
458 std::this_thread::sleep_for(std::chrono::milliseconds(200));
460 EXPECT_EQ(error_q.
Size(), 0);
461 EXPECT_EQ(path_q.
Size(), num_cb_q);
462 EXPECT_EQ(value_q.
Size(), num_cb_q);
475 [&](std::exception_ptr error) { error_q.
PushBack(error); });
482 repo->SendRequest(req).Wait();
485 EXPECT_TRUE(create_q.
Contains(path1));
486 EXPECT_TRUE(create_q.
Contains(path2));
487 EXPECT_EQ(delete_q.
Size(), 0);
488 EXPECT_EQ(error_q.
Size(), 0);
495 repo->SendRequest(req).Wait();
498 EXPECT_EQ(create_q.
Size(), 2);
499 EXPECT_TRUE(delete_q.
Contains(path1));
500 EXPECT_TRUE(delete_q.
Contains(path2));
501 EXPECT_EQ(error_q.
Size(), 0);
510 repo->SendRequest(req).Wait();
513 std::this_thread::sleep_for(std::chrono::milliseconds(200));
515 EXPECT_EQ(create_q.
Size(), 2);
516 EXPECT_EQ(delete_q.
Size(), 2);
517 EXPECT_EQ(error_q.
Size(), 0);
531 [&](std::exception_ptr error) { error_q.
PushBack(error); });
538 repo->SendRequest(req).Wait();
541 EXPECT_TRUE(create_q.
Contains(path1));
542 EXPECT_TRUE(create_q.
Contains(path2));
543 EXPECT_EQ(delete_q.
Size(), 0);
544 EXPECT_EQ(error_q.
Size(), 0);
551 repo->SendRequest(req).Wait();
554 EXPECT_EQ(create_q.
Size(), 2);
555 EXPECT_TRUE(delete_q.
Contains(path1));
556 EXPECT_TRUE(delete_q.
Contains(path2));
557 EXPECT_EQ(error_q.
Size(), 0);
566 repo->SendRequest(req).Wait();
569 std::this_thread::sleep_for(std::chrono::milliseconds(200));
571 EXPECT_EQ(create_q.
Size(), 2);
572 EXPECT_EQ(delete_q.
Size(), 2);
573 EXPECT_EQ(error_q.
Size(), 0);
580 auto sub = repo->Subscribe<
int>(
581 "relative/dp/path"_dppath,
586 [&](std::exception_ptr error) {});
594 auto sub = repo->Subscribe<
int>(
600 [&](std::exception_ptr error) {});
This class provides a wrapper for a data point path.
Definition dataPointPath.hpp:76
An object representing one or more asynchronous I/O requests to a repository.
Definition repositoryIf.hpp:638
void WriteDataPoint(const DataPointPath &path, const T &buffer, std::optional< std::reference_wrapper< MetaData > > metadata=std::nullopt, const CallbackType &callback=nullptr)
Add request to write a datapoint.
Definition repositoryIf.ipp:1557
void DeleteDataPoint(const DataPointPath &path, const CallbackType &callback=nullptr)
Definition repositoryIf.cpp:344
void CreateSymlink(const DataPointPath &dp, const DataPointPath &link, const CallbackType &callback=nullptr)
Definition repositoryIf.cpp:388
void CreateDataPoint(const DataPointPath &path, const T &initial_value, std::optional< std::reference_wrapper< const MetaData > > metadata=std::nullopt, const CallbackType &callback=nullptr)
Add a request to create a new datapoint.
Definition repositoryIf.ipp:1400
void UpdateSymlink(const DataPointPath &dp, const DataPointPath &link, const CallbackType &callback=nullptr)
Definition repositoryIf.cpp:397
Definition repositoryIf.hpp:81
Definition repositoryIf.hpp:77
Definition repositorySubscriberIfTestSuite.hpp:76
void TearDown() override
Definition repositorySubscriberIfTestSuite.hpp:85
DataPointPath path2
Definition repositorySubscriberIfTestSuite.hpp:104
void SetUp() override
Definition repositorySubscriberIfTestSuite.hpp:78
std::shared_ptr< RepositorySubscriberIf > repo
Definition repositorySubscriberIfTestSuite.hpp:102
DataPointPath path1
Definition repositorySubscriberIfTestSuite.hpp:103
DataPointPath link
Definition repositorySubscriberIfTestSuite.hpp:105
Definition repositorySubscriberIfTestSuite.hpp:29
size_t Size()
Definition repositorySubscriberIfTestSuite.hpp:41
void AwaitSize(size_t target_size, std::chrono::seconds timeout=std::chrono::seconds(5))
Definition repositorySubscriberIfTestSuite.hpp:52
void PushBack(const T &path)
Definition repositorySubscriberIfTestSuite.hpp:31
T operator[](size_t idx)
Definition repositorySubscriberIfTestSuite.hpp:36
bool Contains(const T &val)
Definition repositorySubscriberIfTestSuite.hpp:46
Definition fakeClock.cpp:14
std::chrono::milliseconds g_sleep_duration
Definition repositorySubscriberIfTestSuite.hpp:68
TEST_F(Callbacks, CreateDataPointCallback)
Definition repositoryIfTestSuite.hpp:1668
void Sleep()
Definition repositorySubscriberIfTestSuite.hpp:70
Header file for RepositorySubscriberIf and related base classes.