/
smychkov
/
SStorage
Обзор
Документация
Войти
/
smychkov
/
SStorage
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/server/grpc_handler.cpp
189 строк
7 KB
Андрей
feat: top N keys + scan stream (range/top/all) with pagination
14 июл 2026, 12:56
14 июл 2026, 12:56
56c6264
Код
Авторство
О чём код?
#include "grpc_handler.hpp" #include "../core/database.hpp" #include "../util/utils.hpp" #include <iostream> namespace { // На наблюдаемый поток kv не влияет (формат «одно kv на сообщение» сохранён). [[maybe_unused]] constexpr size_t kGrpcScanPageSize = 50; } // namespace namespace sstorage { //============================================================================ // StorageServiceImpl — реализация gRPC-сервиса //============================================================================ class GrpcHandler::StorageServiceImpl final : public Storage::Service { public: explicit StorageServiceImpl(Database& db) : db_(db) {} //------------------------------------------------------------------ // Put //------------------------------------------------------------------ grpc::Status Put(grpc::ServerContext* /*ctx*/, const PutRequest* req, PutResponse* resp) override { bool ok = db_.put(req->key(), req->value()); resp->set_ok(ok); return ok ? grpc::Status::OK : grpc::Status(grpc::INVALID_ARGUMENT, "put failed"); } //------------------------------------------------------------------ // Get //------------------------------------------------------------------ grpc::Status Get(grpc::ServerContext* /*ctx*/, const GetRequest* req, GetResponse* resp) override { auto v = db_.get(req->key()); if (v.has_value()) { resp->set_found(true); resp->set_value(*v); } else { resp->set_found(false); } return grpc::Status::OK; } //------------------------------------------------------------------ // Delete //------------------------------------------------------------------ grpc::Status Delete(grpc::ServerContext* /*ctx*/, const DeleteRequest* req, DeleteResponse* resp) override { bool ok = db_.remove(req->key()); resp->set_ok(ok); return grpc::Status::OK; } //------------------------------------------------------------------ // Scan — streaming (по одной записи за ответ). Контракт S4. //------------------------------------------------------------------ grpc::Status Scan(grpc::ServerContext* /*ctx*/, const ScanRequest* req, grpc::ServerWriter<ScanResponse>* writer) override { auto emit = [writer](const std::string& k, const std::string& v) -> bool { ScanResponse r; r.set_key(k); r.set_value(v); return writer->Write(r); // false = клиент закрыл стрим }; switch (req->mode()) { case SCAN_TOP: { if (req->top_n() == 0) { return grpc::Status::OK; // пустой stream (S4 §3.4) } db_.scanAll(static_cast<size_t>(req->top_n()), emit); return grpc::Status::OK; } case SCAN_ALL: { db_.scanAll(0, emit); return grpc::Status::OK; } case SCAN_RANGE: default: { size_t limit = req->limit() == 0 ? 1000 : req->limit(); auto rows = db_.scan(req->from_key(), req->to_key(), limit); for (const auto& [k, v] : rows) { ScanResponse r; r.set_key(k); r.set_value(v); if (!writer->Write(r)) break; // клиент закрыл стрим } return grpc::Status::OK; } } } //------------------------------------------------------------------ // Stats //------------------------------------------------------------------ grpc::Status Stats(grpc::ServerContext* /*ctx*/, const StatsRequest* /*req*/, StatsResponse* resp) override { auto s = db_.stats(); resp->set_memtable_size(s.memtableSize); resp->set_immutable_count(s.immutableCount); for (auto n : s.levelRecords) resp->add_level_records(n); for (auto n : s.levelFileCounts) resp->add_level_file_counts(n); resp->set_total_records(s.totalRecords); return grpc::Status::OK; } //------------------------------------------------------------------ // Flush / Compact //------------------------------------------------------------------ grpc::Status Flush(grpc::ServerContext* /*ctx*/, const FlushRequest* /*req*/, FlushResponse* resp) override { resp->set_ok(db_.flush()); return grpc::Status::OK; } grpc::Status Compact(grpc::ServerContext* /*ctx*/, const CompactRequest* /*req*/, CompactResponse* resp) override { resp->set_ok(db_.compact()); return grpc::Status::OK; } private: Database& db_; }; //============================================================================ // GrpcHandler — публичный фасад над gRPC-сервисом //============================================================================ GrpcHandler::GrpcHandler(Database& db, uint16_t port) : db_(db), port_(port) {} GrpcHandler::~GrpcHandler() { stop(); } bool GrpcHandler::start() { if (running_) return true; service_ = std::make_unique<StorageServiceImpl>(db_); std::string address = "0.0.0.0:" + std::to_string(port_); grpc::ServerBuilder builder; builder.AddListeningPort(address, grpc::InsecureServerCredentials()); builder.RegisterService(service_.get()); // Лимиты против DoS: // - максимальное число одновременных запросов // - максимальный размер входящего сообщения (ключ + значение + overhead) constexpr int kMaxConcurrentStreams = 128; constexpr int kMaxReceiveMessageSize = 16 * 1024 * 1024; // 16 МБ builder.AddChannelArgument(GRPC_ARG_MAX_CONCURRENT_STREAMS, kMaxConcurrentStreams); builder.SetMaxReceiveMessageSize(kMaxReceiveMessageSize); server_ = builder.BuildAndStart(); if (!server_) { std::cerr << "Failed to start gRPC server on port " << port_ << "\n"; return false; } running_ = true; std::cout << "gRPC server started on port " << port_ << "\n"; return true; } void GrpcHandler::stop() { if (server_) { server_->Shutdown(); server_.reset(); running_ = false; std::cout << "gRPC server stopped\n"; } } bool GrpcHandler::isRunning() const { return running_; } }