/
smychkov
/
SStorage
Обзор
Документация
Войти
/
smychkov
/
SStorage
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/wal/wal.cpp
326 строк
12 KB
Андрей Смычков
security: комплексные фиксы по результатам аудита
25 апр 2026, 10:00
25 апр 2026, 10:00
5f1be8e
Код
Авторство
О чём код?
#include "wal.hpp" #include "../util/crc32c.hpp" #include <cerrno> #include <cstring> #include <dirent.h> #include <fcntl.h> #include <fstream> #include <sys/stat.h> #include <unistd.h> namespace sstorage { namespace { //============================================================================ // Little-endian запись/чтение uint32_t //============================================================================ // Мы явно используем LE-порядок чтобы быть переносимыми между архитектурами. //============================================================================ void writeU32LE(char* out, uint32_t v) { out[0] = static_cast<char>(v & 0xFF); out[1] = static_cast<char>((v >> 8) & 0xFF); out[2] = static_cast<char>((v >> 16) & 0xFF); out[3] = static_cast<char>((v >> 24) & 0xFF); } uint32_t readU32LE(const char* in) { return static_cast<uint8_t>(in[0]) | (static_cast<uint32_t>(static_cast<uint8_t>(in[1])) << 8) | (static_cast<uint32_t>(static_cast<uint8_t>(in[2])) << 16) | (static_cast<uint32_t>(static_cast<uint8_t>(in[3])) << 24); } } //============================================================================ // Конструктор / деструктор //============================================================================ WriteAheadLog::WriteAheadLog(std::string dataDir, uint64_t logNumber, size_t maxQueueBytes) : dataDir_(std::move(dataDir)), logNumber_(logNumber), maxQueueBytes_(maxQueueBytes) { path_ = buildFilename(dataDir_, logNumber_); } WriteAheadLog::~WriteAheadLog() { close(); } //============================================================================ // Открытие файла и запуск writer thread //============================================================================ bool WriteAheadLog::open() { // O_CREAT — создать если нет, O_WRONLY — только запись, // O_APPEND — писать в конец атомарно (важно при restart), // O_NOFOLLOW — не идти по симлинкам (защита от symlink attack) fd_ = ::open(path_.c_str(), O_CREAT | O_WRONLY | O_APPEND | O_NOFOLLOW, 0644); if (fd_ < 0) { return false; } running_ = true; writer_ = std::thread(&WriteAheadLog::writerLoop, this); return true; } //============================================================================ // Закрытие: остановить writer, закрыть fd //============================================================================ void WriteAheadLog::close() { // Уже остановлен? if (!running_.exchange(false)) { if (fd_ >= 0) { ::close(fd_); fd_ = -1; } return; } // Будим ВСЕ ожидающие потоки: и writer, и producer'ов в append() queueCv_.notify_all(); queueDrainCv_.notify_all(); if (writer_.joinable()) { writer_.join(); } if (fd_ >= 0) { ::close(fd_); fd_ = -1; } } //============================================================================ // Формирование фрейма записи //============================================================================ // [crc32c: 4B][length: 4B][payload: N bytes] // CRC считается по (length + payload) — если любой из них испортится, // CRC не сойдётся. //============================================================================ std::string WriteAheadLog::framePayload(const Record& r) { std::string payload; r.serialize(payload); std::string framed; framed.resize(8 + payload.size()); writeU32LE(&framed[4], static_cast<uint32_t>(payload.size())); std::memcpy(&framed[8], payload.data(), payload.size()); // CRC по (length + payload) — первые 4 байта пропускаем uint32_t crc = util::crc32c(&framed[4], 4 + payload.size()); writeU32LE(&framed[0], crc); return framed; } //============================================================================ // Append — блокирующий вызов с ожиданием fsync //============================================================================ bool WriteAheadLog::append(const Record& r) { if (!running_) { return false; } auto pw = std::make_unique<PendingWrite>(); pw->bytes = framePayload(r); auto fut = pw->done.get_future(); size_t payloadSize = pw->bytes.size(); { std::unique_lock<std::mutex> lk(queueMutex_); // Backpressure: ждём пока queueBytes не опустится ниже лимита. // Это защищает от OOM при флуде записями от клиента. // Если writer упадёт (running_ = false) — выходим сразу. if (maxQueueBytes_ > 0) { queueDrainCv_.wait(lk, [&]{ return !running_.load() || queueBytes_ + payloadSize <= maxQueueBytes_; }); if (!running_.load()) return false; } queue_.push(std::move(pw)); queueBytes_ += payloadSize; } queueCv_.notify_one(); // Блокируемся до того, как фоновый поток сделает fsync return fut.get(); } //============================================================================ // Writer loop — group commit //============================================================================ // Собирает все записи из очереди, делает один write()+fsync() на batch, // потом оповещает ожидающих через promise.set_value. // Это даёт amortized O(1) fsync на запись при множестве параллельных клиентов. //============================================================================ void WriteAheadLog::writerLoop() { while (true) { std::vector<std::unique_ptr<PendingWrite>> batch; size_t batchBytes = 0; { std::unique_lock<std::mutex> lk(queueMutex_); queueCv_.wait(lk, [&]{ return !queue_.empty() || !running_; }); // Забираем весь batch одним махом while (!queue_.empty()) { batchBytes += queue_.front()->bytes.size(); batch.push_back(std::move(queue_.front())); queue_.pop(); } // Уменьшаем суммарный размер очереди и будим producer'ов, // которые ждут свободного места (backpressure release). if (batchBytes > 0 && queueBytes_ >= batchBytes) { queueBytes_ -= batchBytes; } else { queueBytes_ = 0; } if (batch.empty() && !running_) { break; } } queueDrainCv_.notify_all(); if (batch.empty()) { continue; } // Склеиваем все записи в один большой буфер — экономим системные вызовы std::string combined; size_t totalSize = 0; for (auto& pw : batch) totalSize += pw->bytes.size(); combined.reserve(totalSize); for (auto& pw : batch) combined.append(pw->bytes); // Один write() + один fsync() на весь batch bool ok = true; size_t written = 0; while (written < combined.size()) { ssize_t n = ::write(fd_, combined.data() + written, combined.size() - written); if (n < 0) { if (errno == EINTR) continue; ok = false; break; } written += static_cast<size_t>(n); } if (ok) { // fsync — гарантирует что данные дошли до диска, а не только в page cache if (::fsync(fd_) != 0) { ok = false; } } // Уведомляем всех писателей пачки о результате for (auto& pw : batch) { pw->done.set_value(ok); } } } //============================================================================ // Replay — чтение всех записей из WAL //============================================================================ // При крэше хвост может быть повреждён (partial write). В таком случае // просто прекращаем чтение на первой же битой записи — это нормальное // поведение, все подтверждённые записи (fsync прошёл) уже прочитаны. //============================================================================ std::vector<Record> WriteAheadLog::replay(const std::string& path) { std::vector<Record> result; std::ifstream file(path, std::ios::binary); if (!file) { return result; } while (file) { char header[8]; file.read(header, 8); if (file.gcount() != 8) { // Не хватило даже на header — корректный EOF или partial write break; } uint32_t expectedCrc = readU32LE(header); uint32_t payloadLen = readU32LE(header + 4); // Защита от невалидно большого payload (подозрение на повреждение) if (payloadLen > 100 * 1024 * 1024) { break; } std::string payload(payloadLen, '\0'); file.read(payload.data(), payloadLen); if (static_cast<size_t>(file.gcount()) != payloadLen) { // Обрыв — partial write, останавливаемся break; } // Проверка CRC: по (length + payload), как при записи в framePayload std::string combined; combined.reserve(4 + payload.size()); combined.append(header + 4, 4); combined.append(payload); uint32_t actualCrc = util::crc32c(combined.data(), combined.size()); if (actualCrc != expectedCrc) { // Повреждённая запись — останавливаемся break; } size_t consumed; auto rec = Record::deserialize(payload.data(), payload.size(), consumed); if (!rec.has_value() || consumed != payload.size()) { break; } result.push_back(std::move(*rec)); } return result; } //============================================================================ // Вспомогательные функции //============================================================================ std::string WriteAheadLog::buildFilename(const std::string& dataDir, uint64_t logNumber) { return dataDir + "/wal_" + std::to_string(logNumber) + ".log"; } std::vector<std::pair<uint64_t, std::string>> WriteAheadLog::listLogFiles(const std::string& dataDir) { std::vector<std::pair<uint64_t, std::string>> result; DIR* dir = ::opendir(dataDir.c_str()); if (!dir) return result; struct dirent* entry; while ((entry = ::readdir(dir)) != nullptr) { std::string name = entry->d_name; // Ожидаем: wal_<number>.log if (name.compare(0, 4, "wal_") != 0) continue; if (name.size() < 9 || name.substr(name.size() - 4) != ".log") continue; std::string numStr = name.substr(4, name.size() - 8); uint64_t num = 0; try { num = std::stoull(numStr); } catch (...) { continue; } result.emplace_back(num, dataDir + "/" + name); } ::closedir(dir); std::sort(result.begin(), result.end(), [](auto& a, auto& b){ return a.first < b.first; }); return result; } bool WriteAheadLog::remove(const std::string& path) { return ::unlink(path.c_str()) == 0; } }