Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- /*
- * Reader-Writer Lock с предотвращением голодания
- * Полная реализация в одном файле
- */
- #include <iostream>
- #include <thread>
- #include <mutex>
- #include <condition_variable>
- #include <queue>
- #include <chrono>
- #include <atomic>
- #include <vector>
- #include <iomanip>
- // Типы запросов в очереди
- enum RequestType {
- READER_REQUEST = 1,
- WRITER_REQUEST = 2
- };
- // Элемент очереди с билетом и типом запроса
- struct QueueItem {
- int ticket;
- RequestType type;
- QueueItem(int t, RequestType rt) : ticket(t), type(rt) {}
- };
- // Основной класс RW Lock с предотвращением голодания
- class RW_Lock {
- private:
- std::mutex mutex_; // Основной мютекс
- std::condition_variable cv_; // Для блокировки потоков
- std::queue<QueueItem> waiting_queue_; // FIFO очередь запросов
- int active_readers_ = 0; // Счетчик активных читателей
- bool active_writer_ = false; // Флаг активного писателя
- int next_ticket_ = 0; // Следующий номер билета
- int current_ticket_ = 0; // Текущий обслуживаемый билет
- public:
- // Захват блокировки для чтения
- void read_lock() {
- std::unique_lock<std::mutex> lock(mutex_);
- // Получаем уникальный билет
- int my_ticket = next_ticket_++;
- // Добавляем себя в очередь
- waiting_queue_.emplace(my_ticket, READER_REQUEST);
- // Ждем своей очереди
- cv_.wait(lock, [this, my_ticket] {
- return current_ticket_ == my_ticket;
- });
- // Ждем пока нет активного писателя
- cv_.wait(lock, [this] {
- return !active_writer_;
- });
- // Увеличиваем счетчик читателей
- ++active_readers_;
- // Убираем себя из очереди и обрабатываем следующих
- waiting_queue_.pop();
- process_next_requests();
- }
- // Освобождение блокировки чтения
- void read_unlock() {
- std::unique_lock<std::mutex> lock(mutex_);
- --active_readers_;
- // Если последний читатель, уведомляем писателей
- if (active_readers_ == 0) {
- cv_.notify_all();
- }
- }
- // Захват блокировки для записи
- void write_lock() {
- std::unique_lock<std::mutex> lock(mutex_);
- // Получаем уникальный билет
- int my_ticket = next_ticket_++;
- // Добавляем себя в очередь
- waiting_queue_.emplace(my_ticket, WRITER_REQUEST);
- // Ждем своей очереди
- cv_.wait(lock, [this, my_ticket] {
- return current_ticket_ == my_ticket;
- });
- // Ждем пока нет активных читателей и писателей
- cv_.wait(lock, [this] {
- return active_readers_ == 0 && !active_writer_;
- });
- // Устанавливаем флаг писателя
- active_writer_ = true;
- // Убираем себя из очереди
- waiting_queue_.pop();
- ++current_ticket_;
- }
- // Освобождение блокировки записи
- void write_unlock() {
- std::unique_lock<std::mutex> lock(mutex_);
- active_writer_ = false;
- process_next_requests();
- cv_.notify_all();
- }
- // Получить статистику (для отладки)
- void print_stats() {
- std::unique_lock<std::mutex> lock(mutex_);
- std::cout << "[DEBUG] Readers: " << active_readers_
- << ", Writer: " << (active_writer_ ? "YES" : "NO")
- << ", Queue size: " << waiting_queue_.size()
- << ", Current ticket: " << current_ticket_ << std::endl;
- }
- private:
- // Обработка следующих запросов в очереди
- void process_next_requests() {
- if (waiting_queue_.empty()) {
- ++current_ticket_;
- } else {
- ++current_ticket_;
- }
- }
- };
- // Класс для тестирования RW Lock
- class RWLockTester {
- private:
- RW_Lock rw_lock_;
- std::atomic<int> shared_data_{0};
- std::atomic<int> reader_ops_{0};
- std::atomic<int> writer_ops_{0};
- std::atomic<bool> test_running_{true};
- public:
- // Тест 1: Базовая функциональность
- void test_basic_functionality() {
- std::cout << "\n=== TEST 1 ===" << std::endl;
- shared_data_ = 0;
- reader_ops_ = 0;
- writer_ops_ = 0;
- const int num_readers = 3;
- const int num_writers = 2;
- const int ops_per_thread = 3;
- std::vector<std::thread> threads;
- // Создаем читателей
- for (int i = 0; i < num_readers; ++i) {
- threads.emplace_back([this, i, ops_per_thread]() {
- for (int j = 0; j < ops_per_thread; ++j) {
- rw_lock_.read_lock();
- int data = shared_data_.load();
- std::this_thread::sleep_for(std::chrono::milliseconds(50));
- std::cout << "Reader " << i+1 << " read: " << data << std::endl;
- reader_ops_++;
- rw_lock_.read_unlock();
- std::this_thread::sleep_for(std::chrono::milliseconds(30));
- }
- });
- }
- // Создаем писателей
- for (int i = 0; i < num_writers; ++i) {
- threads.emplace_back([this, i, ops_per_thread]() {
- for (int j = 0; j < ops_per_thread; ++j) {
- rw_lock_.write_lock();
- int old_val = shared_data_.load();
- std::this_thread::sleep_for(std::chrono::milliseconds(80));
- shared_data_.store(old_val + 1);
- std::cout << "Writer " << i+1 << " changes: "
- << old_val << " -> " << shared_data_.load() << std::endl;
- writer_ops_++;
- rw_lock_.write_unlock();
- std::this_thread::sleep_for(std::chrono::milliseconds(40));
- }
- });
- }
- // Ждем завершения
- for (auto& t : threads) {
- t.join();
- }
- // Проверяем результаты
- std::cout << "\nResults 1:" << std::endl;
- std::cout << "Final: " << shared_data_.load() << std::endl;
- std::cout << "Readings: " << reader_ops_.load() << std::endl;
- std::cout << "Changes: " << writer_ops_.load() << std::endl;
- std::cout << "Waitings: " << num_writers * ops_per_thread << std::endl;
- bool passed = (writer_ops_.load() == num_writers * ops_per_thread) &&
- (shared_data_.load() == writer_ops_.load());
- std::cout << "Тест 1: " << (passed ? "OK" : "BAD") << std::endl;
- }
- // Тест 2: Предотвращение голодания
- void test_starvation_prevention() {
- std::cout << "\n=== TEST 2 ===" << std::endl;
- shared_data_ = 0;
- reader_ops_ = 0;
- writer_ops_ = 0;
- test_running_ = true;
- std::vector<std::chrono::steady_clock::time_point> writer_start_times(2);
- std::vector<std::chrono::steady_clock::time_point> writer_end_times(2);
- // Создаем много непрерывных читателей
- std::vector<std::thread> reader_threads;
- for (int i = 0; i < 5; ++i) {
- reader_threads.emplace_back([this, i]() {
- int local_ops = 0;
- while (test_running_.load() && local_ops < 20) {
- rw_lock_.read_lock();
- int data = shared_data_.load();
- std::this_thread::sleep_for(std::chrono::milliseconds(20));
- reader_ops_++;
- local_ops++;
- rw_lock_.read_unlock();
- std::this_thread::sleep_for(std::chrono::milliseconds(10));
- }
- });
- }
- // Создаем писателей с измерением времени ожидания
- std::vector<std::thread> writer_threads;
- for (int i = 0; i < 2; ++i) {
- writer_threads.emplace_back([this, i, &writer_start_times, &writer_end_times]() {
- std::this_thread::sleep_for(std::chrono::milliseconds(100 + i * 300));
- writer_start_times[i] = std::chrono::steady_clock::now();
- std::cout << "Writer " << i+1 << " requests..." << std::endl;
- rw_lock_.write_lock();
- writer_end_times[i] = std::chrono::steady_clock::now();
- int old_val = shared_data_.load();
- shared_data_.store(old_val + 10);
- writer_ops_++;
- auto wait_ms = std::chrono::duration_cast<std::chrono::milliseconds>
- (writer_end_times[i] - writer_start_times[i]).count();
- std::cout << "Writer " << i+1 << " got " << wait_ms
- << "changes: " << old_val << " -> " << shared_data_.load() << std::endl;
- std::this_thread::sleep_for(std::chrono::milliseconds(50));
- rw_lock_.write_unlock();
- });
- }
- // Запускаем тест на 1.5 секунды
- std::this_thread::sleep_for(std::chrono::milliseconds(1500));
- test_running_ = false;
- // Ждем завершения
- for (auto& t : reader_threads) {
- t.join();
- }
- for (auto& t : writer_threads) {
- t.join();
- }
- // Анализируем результаты
- std::cout << "\Results 2:" << std::endl;
- std::cout << "Readings: " << reader_ops_.load() << std::endl;
- std::cout << "Changes: " << writer_ops_.load() << std::endl;
- std::cout << "Final: " << shared_data_.load() << std::endl;
- std::cout << "Waiting time:" << std::endl;
- for (int i = 0; i < 2; ++i) {
- if (writer_end_times[i] > writer_start_times[i]) {
- auto wait_ms = std::chrono::duration_cast<std::chrono::milliseconds>
- (writer_end_times[i] - writer_start_times[i]).count();
- std::cout << " Writer " << i+1 << ": " << wait_ms << "мс" << std::endl;
- }
- }
- bool no_starvation = (writer_ops_.load() == 2) && (shared_data_.load() == 20);
- std::cout << "Starving: " << (no_starvation ? "OK" : "BAD") << std::endl;
- }
- // Тест 3: Параллельные читатели
- void test_concurrent_readers() {
- std::cout << "\n=== TEST 3 ===" << std::endl;
- shared_data_ = 42;
- reader_ops_ = 0;
- const int num_readers = 6;
- std::vector<std::thread> threads;
- auto start_time = std::chrono::steady_clock::now();
- // Создаем читателей которые должны работать параллельно
- for (int i = 0; i < num_readers; ++i) {
- threads.emplace_back([this, i]() {
- rw_lock_.read_lock();
- auto thread_start = std::chrono::steady_clock::now();
- std::cout << "Reader " << i+1 << " starts" << std::endl;
- // Имитируем долгое чтение
- std::this_thread::sleep_for(std::chrono::milliseconds(200));
- int data = shared_data_.load();
- auto thread_end = std::chrono::steady_clock::now();
- auto duration = std::chrono::duration_cast<std::chrono::milliseconds>
- (thread_end - thread_start).count();
- std::cout << "Reader " << i+1 << " ends "
- << duration << "data: " << data << std::endl;
- reader_ops_++;
- rw_lock_.read_unlock();
- });
- }
- for (auto& t : threads) {
- t.join();
- }
- auto end_time = std::chrono::steady_clock::now();
- auto total_time = std::chrono::duration_cast<std::chrono::milliseconds>
- (end_time - start_time).count();
- std::cout << "\nResults 3:" << std::endl;
- std::cout << "Time: " << total_time << "мс" << std::endl;
- std::cout << "Readings: " << reader_ops_.load() << std::endl;
- bool concurrent = total_time < 400; // Даем запас
- std::cout << "Reading: " << (concurrent ? "OK" : "BAD") << std::endl;
- }
- // Запуск всех тестов
- void run_all_tests() {
- std::cout << "=== READER-WRITER LOCK ===" << std::endl;
- std::cout << "Ticket System + FIFO Queue" << std::endl;
- test_basic_functionality();
- std::this_thread::sleep_for(std::chrono::milliseconds(500));
- test_concurrent_readers();
- std::this_thread::sleep_for(std::chrono::milliseconds(500));
- test_starvation_prevention();
- }
- };
- // Основное тело
- int main() {
- try {
- RWLockTester tester;
- tester.run_all_tests();
- std::cout << "✓ Ticket System" << std::endl;
- std::cout << "✓ FIFO Queue" << std::endl;
- std::cout << "✓ Starving System" << std::endl;
- std::cout << "✓ std::condition_variable without busy waiting" << std::endl;
- std::cout << "✓ RW_Lock Class" << std::endl;
- std::cout << "✓ Thread-safe" << std::endl;
- } catch (const std::exception& e) {
- std::cerr << "Error: " << e.what() << std::endl;
- return 1;
- }
- return 0;
- }
Advertisement
Add Comment
Please, Sign In to add comment