13#ifndef RTCTK_COMPONENTFRAMEWORK_TEST_REPOSITORYSUBSCRIBERIFTESTSUITE_HPP
14#define RTCTK_COMPONENTFRAMEWORK_TEST_REPOSITORYSUBSCRIBERIFTESTSUITE_HPP
19#include <gtest/gtest.h>
27static std::shared_ptr<RepositorySubscriberIf> MakeRepository();
33 std::lock_guard lock{m_mutex};
34 m_data.push_back(path);
38 std::lock_guard lock{m_mutex};
39 return m_data.at(idx);
43 std::lock_guard lock{m_mutex};
48 std::lock_guard lock{m_mutex};
49 auto it = std::find(m_data.begin(), m_data.end(), val);
50 return it != m_data.end();
53 void AwaitSize(
size_t target_size, std::chrono::seconds timeout = std::chrono::seconds(5)) {
54 auto t_start = std::chrono::steady_clock::now();
55 while (
Size() < target_size) {
56 std::this_thread::sleep_for(std::chrono::microseconds(100));
57 if ((std::chrono::steady_clock::now() - t_start) > timeout) {
58 FAIL() <<
"ThreadSafeQ::AwaitSize: Timed out waiting for size "
59 << std::to_string(target_size);
66 std::vector<T> m_data;
80 repo = MakeRepository();
81 path1 =
"/foo"_dppath;
82 path2 =
"/bar"_dppath;
83 link =
"/link"_dppath;
87 if (
repo ==
nullptr) {
103 std::shared_ptr<RepositorySubscriberIf>
repo;
115 repo->SendRequest(req).Wait();
125 auto sub1 = repo->Subscribe<
int>(
134 [&](std::exception_ptr error) { error_q_1.
PushBack(error); });
136 auto sub2 = repo->Subscribe(
151 repo->SendRequest(req).Wait();
156 EXPECT_TRUE(path_q_1.
Contains(path1));
157 EXPECT_TRUE(value_q_1.
Contains(int_write));
158 EXPECT_TRUE(path_q_2.
Contains(path2));
159 EXPECT_EQ(error_q_1.
Size(), 0);
165 num_cb_q_1 = path_q_1.
Size();
166 num_cb_q_2 = path_q_2.
Size();
173 repo->SendRequest(req).Wait();
176 std::this_thread::sleep_for(std::chrono::milliseconds(200));
179 EXPECT_EQ(error_q_1.
Size(), 0);
180 EXPECT_EQ(path_q_1.
Size(), num_cb_q_1);
181 EXPECT_EQ(value_q_1.
Size(), num_cb_q_1);
182 EXPECT_EQ(path_q_2.
Size(), num_cb_q_2);
187 repo->CreateDataPoint(path1, 0);
194 auto sub1 = repo->Subscribe(
202 auto sub2 = repo->Subscribe(
213 repo->WriteDataPoint(path1, 3);
217 EXPECT_TRUE(path_q_1.
Contains(path1));
218 EXPECT_TRUE(path_q_2.
Contains(path1));
222 num_cb_q_1 = path_q_1.
Size();
223 num_cb_q_2 = path_q_2.
Size();
226 repo->WriteDataPoint(path1, 4);
231 std::this_thread::sleep_for(std::chrono::milliseconds(200));
233 EXPECT_EQ(path_q_1.
Size(), (num_cb_q_1 + 1));
234 EXPECT_EQ(path_q_2.
Size(), num_cb_q_2);
243 repo->SendRequest(req).Wait();
246 std::atomic_bool stop =
false;
249 auto writer_func = [&] {
251 while (stop ==
false) {
252 repo->WriteDataPoint(path1, value);
258 auto subscriber_func = [&] {
259 for (
unsigned i = 0; i < 1000; i++) {
262 auto sub = repo->Subscribe(
276 std::thread t1{writer_func};
277 std::thread t2{subscriber_func};
288 repo->SendRequest(req).Wait();
296 auto sub = repo->Subscribe<
int>(
305 [&](std::exception_ptr error) { error_q.
PushBack(error); });
313 repo->WriteDataPoint(path1, int_write);
318 EXPECT_TRUE(value_q.
Contains(int_write));
319 EXPECT_EQ(error_q.
Size(), 0);
325 repo->WriteDataPoint(link, int_write);
330 EXPECT_TRUE(value_q.
Contains(int_write));
331 EXPECT_EQ(error_q.
Size(), 0);
338 int num_cb_q = path_q.
Size();
341 repo->WriteDataPoint(path1, int_write);
344 std::this_thread::sleep_for(std::chrono::milliseconds(200));
346 EXPECT_EQ(error_q.
Size(), 0);
347 EXPECT_EQ(path_q.
Size(), num_cb_q);
348 EXPECT_EQ(value_q.
Size(), num_cb_q);
358 repo->SendRequest(req).Wait();
365 auto sub = repo->Subscribe<
int>(
374 [&](std::exception_ptr error) { error_q.
PushBack(error); });
381 repo->WriteDataPoint(path1, int_write);
386 EXPECT_TRUE(value_q.
Contains(int_write));
387 EXPECT_EQ(error_q.
Size(), 0);
390 int num_cb_q = path_q.
Size();
391 repo->DeleteDataPoint(path1);
394 std::this_thread::sleep_for(std::chrono::milliseconds(200));
396 EXPECT_EQ(error_q.
Size(), 0);
397 EXPECT_EQ(path_q.
Size(), num_cb_q);
398 EXPECT_EQ(value_q.
Size(), num_cb_q);
404 int int_create_2 = 100;
409 repo->SendRequest(req).Wait();
416 auto sub = repo->Subscribe<
int>(
425 [&](std::exception_ptr error) { error_q.
PushBack(error); });
432 repo->WriteDataPoint(path1, int_write);
437 EXPECT_TRUE(value_q.
Contains(int_write));
438 EXPECT_EQ(error_q.
Size(), 0);
444 repo->SendRequest(req).Wait();
451 EXPECT_EQ(error_q.
Size(), 0);
454 int num_cb_q = path_q.
Size();
458 repo->SendRequest(req).Wait();
461 repo->WriteDataPoint(path1, int_write);
462 repo->WriteDataPoint(path2, int_write);
465 std::this_thread::sleep_for(std::chrono::milliseconds(200));
467 EXPECT_EQ(error_q.
Size(), 0);
468 EXPECT_EQ(path_q.
Size(), num_cb_q);
469 EXPECT_EQ(value_q.
Size(), num_cb_q);
482 [&](std::exception_ptr error) { error_q.
PushBack(error); });
489 repo->SendRequest(req).Wait();
492 EXPECT_TRUE(create_q.
Contains(path1));
493 EXPECT_TRUE(create_q.
Contains(path2));
494 EXPECT_EQ(delete_q.
Size(), 0);
495 EXPECT_EQ(error_q.
Size(), 0);
502 repo->SendRequest(req).Wait();
505 EXPECT_EQ(create_q.
Size(), 2);
506 EXPECT_TRUE(delete_q.
Contains(path1));
507 EXPECT_TRUE(delete_q.
Contains(path2));
508 EXPECT_EQ(error_q.
Size(), 0);
517 repo->SendRequest(req).Wait();
520 std::this_thread::sleep_for(std::chrono::milliseconds(200));
522 EXPECT_EQ(create_q.
Size(), 2);
523 EXPECT_EQ(delete_q.
Size(), 2);
524 EXPECT_EQ(error_q.
Size(), 0);
538 [&](std::exception_ptr error) { error_q.
PushBack(error); });
545 repo->SendRequest(req).Wait();
548 EXPECT_TRUE(create_q.
Contains(path1));
549 EXPECT_TRUE(create_q.
Contains(path2));
550 EXPECT_EQ(delete_q.
Size(), 0);
551 EXPECT_EQ(error_q.
Size(), 0);
558 repo->SendRequest(req).Wait();
561 EXPECT_EQ(create_q.
Size(), 2);
562 EXPECT_TRUE(delete_q.
Contains(path1));
563 EXPECT_TRUE(delete_q.
Contains(path2));
564 EXPECT_EQ(error_q.
Size(), 0);
573 repo->SendRequest(req).Wait();
576 std::this_thread::sleep_for(std::chrono::milliseconds(200));
578 EXPECT_EQ(create_q.
Size(), 2);
579 EXPECT_EQ(delete_q.
Size(), 2);
580 EXPECT_EQ(error_q.
Size(), 0);
587 auto sub = repo->Subscribe<
int>(
588 "relative/dp/path"_dppath,
593 [&](std::exception_ptr error) {});
601 auto sub = repo->Subscribe<
int>(
607 [&](std::exception_ptr error) {});
This class provides a wrapper for a data point path.
Definition dataPointPath.hpp:77
An object representing one or more asynchronous I/O requests to a repository.
Definition repositoryIf.hpp:634
void WriteDataPoint(const DataPointPath &path, const T &buffer, std::optional< std::reference_wrapper< MetaData > > metadata=std::nullopt, const CallbackType &callback=nullptr)
Definition repositoryIf.ipp:1538
void DeleteDataPoint(const DataPointPath &path, const CallbackType &callback=nullptr)
Definition repositoryIf.cpp:343
void CreateSymlink(const DataPointPath &dp, const DataPointPath &link, const CallbackType &callback=nullptr)
Definition repositoryIf.cpp:383
void CreateDataPoint(const DataPointPath &path, const T &initial_value, std::optional< std::reference_wrapper< const MetaData > > metadata=std::nullopt, const CallbackType &callback=nullptr)
Definition repositoryIf.ipp:1385
void UpdateSymlink(const DataPointPath &dp, const DataPointPath &link, const CallbackType &callback=nullptr)
Definition repositoryIf.cpp:392
Definition repositoryIf.hpp:82
Definition repositoryIf.hpp:78
Definition repositorySubscriberIfTestSuite.hpp:77
void TearDown() override
Definition repositorySubscriberIfTestSuite.hpp:86
DataPointPath path2
Definition repositorySubscriberIfTestSuite.hpp:105
void SetUp() override
Definition repositorySubscriberIfTestSuite.hpp:79
std::shared_ptr< RepositorySubscriberIf > repo
Definition repositorySubscriberIfTestSuite.hpp:103
DataPointPath path1
Definition repositorySubscriberIfTestSuite.hpp:104
DataPointPath link
Definition repositorySubscriberIfTestSuite.hpp:106
Definition repositorySubscriberIfTestSuite.hpp:30
size_t Size()
Definition repositorySubscriberIfTestSuite.hpp:42
void AwaitSize(size_t target_size, std::chrono::seconds timeout=std::chrono::seconds(5))
Definition repositorySubscriberIfTestSuite.hpp:53
void PushBack(const T &path)
Definition repositorySubscriberIfTestSuite.hpp:32
T operator[](size_t idx)
Definition repositorySubscriberIfTestSuite.hpp:37
bool Contains(const T &val)
Definition repositorySubscriberIfTestSuite.hpp:47
Definition fakeClock.cpp:15
std::chrono::milliseconds g_sleep_duration
Definition repositorySubscriberIfTestSuite.hpp:69
TEST_F(Callbacks, CreateDataPointCallback)
Definition repositoryIfTestSuite.hpp:1662
void Sleep()
Definition repositorySubscriberIfTestSuite.hpp:71
Header file for RepositorySubscriberIf and related base classes.