summaryrefslogtreecommitdiff
path: root/libs/tdlib/td/tdutils/test/MpmcWaiter.cpp
blob: e27e2177139fb5f46092ecd6d5f53b99e587278d (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
//
// Copyright Aliaksei Levin (levlam@telegram.org), Arseny Smirnov (arseny30@gmail.com) 2014-2018
//
// Distributed under the Boost Software License, Version 1.0. (See accompanying
// file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
//
#include "td/utils/MpmcWaiter.h"
#include "td/utils/port/sleep.h"
#include "td/utils/port/thread.h"
#include "td/utils/Random.h"
#include "td/utils/tests.h"

#include <atomic>

#if !TD_THREAD_UNSUPPORTED
TEST(MpmcWaiter, stress_one_one) {
  td::Stage run;
  td::Stage check;

  std::vector<td::thread> threads;
  std::atomic<size_t> value;
  size_t write_cnt = 10;
  std::unique_ptr<td::MpmcWaiter> waiter;
  size_t threads_n = 2;
  for (size_t i = 0; i < threads_n; i++) {
    threads.push_back(td::thread([&, id = static_cast<td::uint32>(i)] {
      for (td::uint64 round = 1; round < 100000; round++) {
        if (id == 0) {
          value = 0;
          waiter = std::make_unique<td::MpmcWaiter>();
          write_cnt = td::Random::fast(1, 10);
        }
        run.wait(round * threads_n);
        if (id == 1) {
          for (size_t i = 0; i < write_cnt; i++) {
            value.store(i + 1, std::memory_order_relaxed);
            waiter->notify();
          }
        } else {
          int yields = 0;
          for (size_t i = 1; i <= write_cnt; i++) {
            while (true) {
              auto x = value.load(std::memory_order_relaxed);
              if (x >= i) {
                break;
              }
              yields = waiter->wait(yields, id);
            }
            yields = waiter->stop_wait(yields, id);
          }
        }
        check.wait(round * threads_n);
      }
    }));
  }
  for (auto &thread : threads) {
    thread.join();
  }
}
TEST(MpmcWaiter, stress) {
  td::Stage run;
  td::Stage check;

  std::vector<td::thread> threads;
  size_t write_n;
  size_t read_n;
  std::atomic<size_t> write_pos;
  std::atomic<size_t> read_pos;
  size_t end_pos;
  size_t write_cnt;
  size_t threads_n = 20;
  std::unique_ptr<td::MpmcWaiter> waiter;
  for (size_t i = 0; i < threads_n; i++) {
    threads.push_back(td::thread([&, id = static_cast<td::uint32>(i)] {
      for (td::uint64 round = 1; round < 1000; round++) {
        if (id == 0) {
          write_n = td::Random::fast(1, 10);
          read_n = td::Random::fast(1, 10);
          write_cnt = td::Random::fast(1, 50);
          end_pos = write_n * write_cnt;
          write_pos = 0;
          read_pos = 0;
          waiter = std::make_unique<td::MpmcWaiter>();
        }
        run.wait(round * threads_n);
        if (id <= write_n) {
          for (size_t i = 0; i < write_cnt; i++) {
            if (td::Random::fast(0, 20) == 0) {
              td::usleep_for(td::Random::fast(1, 300));
            }
            write_pos.fetch_add(1, std::memory_order_relaxed);
            waiter->notify();
          }
        } else if (id > 10 && id - 10 <= read_n) {
          int yields = 0;
          while (true) {
            auto x = read_pos.load(std::memory_order_relaxed);
            if (x == end_pos) {
              break;
            }
            if (x == write_pos.load(std::memory_order_relaxed)) {
              yields = waiter->wait(yields, id);
              continue;
            }
            yields = waiter->stop_wait(yields, id);
            read_pos.compare_exchange_strong(x, x + 1, std::memory_order_relaxed);
          }
        }
        check.wait(round * threads_n);
      }
    }));
  }
  for (auto &thread : threads) {
    thread.join();
  }
}
#endif  // !TD_THREAD_UNSUPPORTED