/
isachkov
/
crypto_server
Обзор
Документация
Войти
/
isachkov
/
crypto_server
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
master
src/executor.cpp
166 строк
4 KB
Ivan Sachkov
CryptoServer
15 май 2024, 20:26
15 май 2024, 20:26
9de0f1f
Код
Авторство
О чём код?
#include <crypto_chat/executor.hpp> #include <fcntl.h> #include <unistd.h> #include <array> #include <cstring> #include <iostream> bool Executor::init() { epoll_fd_ = epoll_create1(0); if (epoll_fd_ == -1) { return false; } return true; } bool Executor::on_read(Descriptor desc, OnRead handler) { const auto is_added = add(desc, EPOLLIN | EPOLLET); if (!is_added) { return is_added; } read_handlers_.emplace(desc, std::move(handler)); return true; } bool Executor::on_listen(Descriptor desc, OnListen handler) { const auto is_added = add(desc, EPOLLIN | EPOLLET); if (!is_added) { return is_added; } listen_fd_ = desc; on_listen_ = std::move(handler); return true; } void Executor::on_disconnect(OnDisconnect handler) { on_disconnect_ = std::move(handler); } void Executor::run() { while (running_) { constexpr auto event_pool_size = 32; std::array<epoll_event, event_pool_size> events; const auto event_count = epoll_wait(epoll_fd_, events.data(), event_pool_size, -1); if (event_count == -1) { running_ = false; return; } for (auto i = 0; i < event_count; ++i) { auto& event = events[i]; if(is_listener_event(event)) { process_listen(event); } else if (event.events & EPOLLIN) { process_read(event); } } } } void Executor::stop() { running_ = false; } void Executor::process_listen(epoll_event& event) { if (!on_listen_) { std::cerr << "No listen handler is set!\n"; return; } on_listen_(); } bool Executor::add(Descriptor descriptor, EventFlags flags) { if (!epoll_fd_) { return false; } // Force all descriptors to use nonblocking io if (fcntl(descriptor, F_SETFL, O_NONBLOCK) == -1) { return false; } epoll_event desc_event; desc_event.data.fd = descriptor; desc_event.events = flags; const auto result = epoll_ctl(epoll_fd_, EPOLL_CTL_ADD, descriptor, &desc_event); if (result != 0) { return false; } return true; } void Executor::remove(Descriptor descriptor) { if (!epoll_fd_) { return; } epoll_ctl(epoll_fd_, EPOLL_CTL_DEL, descriptor, nullptr); close(descriptor); } void Executor::process_read(epoll_event& event) { constexpr auto read_buffer_size = 4096; // In a highly parallel version one should implement some thread-local // storages, yet we are bound to a single threaded solution here static std::array<uint8_t, read_buffer_size> read_buffer; // No handler - no fun auto iter = read_handlers_.find(event.data.fd); if (iter == read_handlers_.end()) { return; } auto read_count = 0; // In a production grade system this has to be pooled std::vector<uint8_t> output; // Do at least once: do { // 1. Read the data, we are nonblocking so expect -1 and EAGAIN on end read_count = read(event.data.fd, read_buffer.data(), read_buffer_size); // 2. Read nothing -> underlying connection/medium is dead if(!read_count) { break; } // 3. May pre-reserve it, but for illustration project just copy as-is std::copy(read_buffer.begin(), read_buffer.begin() + read_count, std::back_inserter(output)); } while (read_count > 0); // Got and error and it's not EAGAIN? Do nothing, or add error handling here. if (read_count < 0 && errno != EAGAIN) { std::cerr << "Read error: " << strerror(errno) << '\n'; return; } // Empty output means connection is closed/file deleted if (!output.size()) { if (on_disconnect_) { on_disconnect_(event.data.fd); } remove(event.data.fd); return; } iter->second(std::move(output)); } bool Executor::is_listener_event(epoll_event& event) { return listen_fd_.has_value() && *listen_fd_ == event.data.fd; }