/
smychkov
/
SStorage
Обзор
Документация
Войти
/
smychkov
/
SStorage
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/server/http_handler.cpp
534 строки
22 KB
Андрей
feat: top N keys + scan stream (range/top/all) with pagination
14 июл 2026, 12:56
14 июл 2026, 12:56
56c6264
Код
Авторство
О чём код?
#include "http_handler.hpp" #include "../core/database.hpp" #include "../util/utils.hpp" #include <algorithm> #include <cstring> #include <iostream> #include <sstream> #include <stdexcept> #include <string> #include <utility> #include <vector> namespace sstorage { namespace { //============================================================================ // Escape для JSON-строк //============================================================================ // Мы храним произвольные байты (bytes key/value как в Bigtable), поэтому // значение может содержать невалидный UTF-8. JSON требует валидного UTF-8 // или \uXXXX escape. Чтобы не отправлять битую кодировку клиенту: // - ASCII печатные (0x20..0x7E) — как есть, кроме '"' и '\\' // - управляющие (<0x20) — \uXXXX // - всё >= 0x80 — \uXXXX (гарантированно валидный JSON независимо // от того, валидный ли UTF-8 был в исходных байтах) // Это защищает от content smuggling через намеренно повреждённый UTF-8. //============================================================================ std::string jsonEscape(const std::string& s) { std::string out; out.reserve(s.size() + 8); for (unsigned char uc : s) { char c = static_cast<char>(uc); switch (c) { case '"': out += "\\\""; break; case '\\': out += "\\\\"; break; case '\n': out += "\\n"; break; case '\r': out += "\\r"; break; case '\t': out += "\\t"; break; case '\b': out += "\\b"; break; case '\f': out += "\\f"; break; default: if (uc < 0x20 || uc >= 0x80) { // Все управляющие и non-ASCII байты — в \uXXXX. // Это даёт гарантированно валидный JSON независимо // от содержимого исходной строки. char buf[8]; std::snprintf(buf, sizeof(buf), "\\u%04x", static_cast<unsigned int>(uc)); out += buf; } else { out += c; } } } return out; } //============================================================================ // URL-decode для значений параметров //============================================================================ std::string urlDecode(const std::string& s) { std::string out; out.reserve(s.size()); for (size_t i = 0; i < s.size(); ++i) { if (s[i] == '%' && i + 2 < s.size()) { char hex[3] = {s[i+1], s[i+2], '\0'}; out.push_back(static_cast<char>(std::strtol(hex, nullptr, 16))); i += 2; } else if (s[i] == '+') { out.push_back(' '); } else { out.push_back(s[i]); } } return out; } //============================================================================ // Константы cursor-пагинации /scan. См. contract S3 §8. //============================================================================ constexpr size_t kHttpPageSize = 50; constexpr size_t kHttpPageSizeMax = 1000; //============================================================================ // hex-кодировка opaque-курсора (§4.2). Байты → lowercase hex. //============================================================================ std::string hexEncode(const std::string& raw) { static const char* kHex = "0123456789abcdef"; std::string out; out.reserve(raw.size() * 2); for (unsigned char uc : raw) { out.push_back(kHex[uc >> 4]); out.push_back(kHex[uc & 0x0F]); } return out; } //============================================================================ // hex-декодировка курсора. Нечётная длина / не-hex символ → "" (§4.2, §6). //============================================================================ std::string hexDecode(const std::string& hex) { if (hex.size() % 2 != 0) return ""; auto nibble = [](char c) -> int { if (c >= '0' && c <= '9') return c - '0'; if (c >= 'a' && c <= 'f') return c - 'a' + 10; if (c >= 'A' && c <= 'F') return c - 'A' + 10; return -1; }; std::string out; out.reserve(hex.size() / 2); for (size_t i = 0; i < hex.size(); i += 2) { int hi = nibble(hex[i]); int lo = nibble(hex[i + 1]); if (hi < 0 || lo < 0) return ""; out.push_back(static_cast<char>((hi << 4) | lo)); } return out; } //============================================================================ // Парсинг неотрицательного целого из query-параметра (§3.1). // Пустое / нечисло / отрицательное / с лишними символами → 0. //============================================================================ size_t parseUIntParam(const std::string& s) { if (s.empty()) return 0; for (char c : s) { if (c < '0' || c > '9') return 0; } try { unsigned long long v = std::stoull(s); return static_cast<size_t>(v); } catch (...) { return 0; } } } //============================================================================ // Конструктор / деструктор //============================================================================ HttpHandler::HttpHandler(Database& db, uint16_t port) : db_(db), port_(port) {} HttpHandler::~HttpHandler() { stop(); } //============================================================================ // start / stop //============================================================================ bool HttpHandler::start() { if (running_) return true; // Лимиты для защиты от DoS: // - максимум 128 одновременных соединений // - максимум 16 на один IP // - таймаут неактивного соединения 30 секунд // - максимальный размер тела запроса 16 МБ constexpr unsigned int kMaxConnections = 128; constexpr unsigned int kMaxConnectionsPerIp = 16; constexpr unsigned int kConnectionTimeoutSec = 30; constexpr size_t kMaxBodySize = 16 * 1024 * 1024; daemon_ = MHD_start_daemon( MHD_USE_INTERNAL_POLLING_THREAD, port_, nullptr, nullptr, &HttpHandler::handleRequest, this, MHD_OPTION_CONNECTION_LIMIT, kMaxConnections, MHD_OPTION_PER_IP_CONNECTION_LIMIT, kMaxConnectionsPerIp, MHD_OPTION_CONNECTION_TIMEOUT, kConnectionTimeoutSec, MHD_OPTION_CONNECTION_MEMORY_LIMIT, kMaxBodySize, MHD_OPTION_END); if (!daemon_) { std::cerr << "Failed to start HTTP server on port " << port_ << "\n"; return false; } running_ = true; std::cout << "HTTP server started on port " << port_ << "\n"; return true; } void HttpHandler::stop() { if (daemon_) { MHD_stop_daemon(daemon_); daemon_ = nullptr; running_ = false; std::cout << "HTTP server stopped\n"; } } bool HttpHandler::isRunning() const { return running_; } //============================================================================ // Главный callback MHD //============================================================================ // Вызывается несколько раз на один запрос: // 1. Первый вызов — тело ещё не пришло (только заголовки). Сохраняем state. // 2. Последующие — передают куски body (POST/PUT). Накапливаем. // 3. Финальный — upload_data_size = 0. Обрабатываем запрос. //============================================================================ MHD_Result HttpHandler::handleRequest(void* cls, MHD_Connection* connection, const char* url, const char* method, const char* /*version*/, const char* upload_data, size_t* upload_data_size, void** con_cls) { HttpHandler* self = static_cast<HttpHandler*>(cls); // Первый вызов — создаём state if (*con_cls == nullptr) { auto* state = new ConnectionState(); *con_cls = state; return MHD_YES; } auto* state = static_cast<ConnectionState*>(*con_cls); // Накопление body (может прийти по частям) if (*upload_data_size > 0) { state->body.append(upload_data, *upload_data_size); *upload_data_size = 0; return MHD_YES; } // Финальный вызов — обрабатываем std::string urlStr(url); MHD_Result result = MHD_NO; if (std::strcmp(method, "GET") == 0) { result = self->processGet(connection, urlStr); } else if (std::strcmp(method, "PUT") == 0) { result = self->processPut(connection, urlStr, state->body); } else if (std::strcmp(method, "DELETE") == 0) { result = self->processDelete(connection, urlStr); } else if (std::strcmp(method, "POST") == 0) { result = self->processPost(connection, urlStr); } else { result = self->sendError(connection, "Method not allowed", 405); } delete state; *con_cls = nullptr; return result; } //============================================================================ // GET //============================================================================ MHD_Result HttpHandler::processGet(MHD_Connection* connection, const std::string& url) { // Путь без query std::string path = url; auto q = path.find('?'); if (q != std::string::npos) path = path.substr(0, q); if (path == "/health") { return sendJson(connection, R"({"status":"ok"})", 200); } // GET /kv/{key} if (path.compare(0, 4, "/kv/") == 0) { std::string key = urlDecode(path.substr(4)); auto v = db_.get(key); if (v.has_value()) { return sendBinary(connection, *v, 200); } return sendError(connection, "Not found", 404); } // GET /scan — режимы top | диапазон | все (cursor-пагинация). См. contract S3. if (path == "/scan") { return handleScan(connection, url); } // GET /stats if (path == "/stats") { auto s = db_.stats(); std::ostringstream oss; oss << "{\"memtable_size\":" << s.memtableSize << ",\"immutable_count\":" << s.immutableCount << ",\"total_records\":" << s.totalRecords << ",\"levels\":["; for (size_t i = 0; i < s.levelFileCounts.size(); ++i) { if (i) oss << ","; oss << "{\"level\":" << i << ",\"files\":" << s.levelFileCounts[i] << ",\"records\":" << s.levelRecords[i] << "}"; } oss << "]}"; return sendJson(connection, oss.str(), 200); } return sendError(connection, "Not found", 404); } //============================================================================ // GET /scan — режимы top N | диапазон | все (cursor-пагинация). См. contract S3. //============================================================================ MHD_Result HttpHandler::handleScan(MHD_Connection* connection, const std::string& /*url*/) { // libmicrohttpd передаёт в callback url БЕЗ query-string; параметры // разбираются через MHD_GET_ARGUMENT_KIND (уже URL-decoded). auto param = [&](const char* name) -> const char* { return MHD_lookup_connection_value(connection, MHD_GET_ARGUMENT_KIND, name); }; // Курсор → lastKey (hex-decoded). Битый/невалидный → "" (с начала). См. §4.2/§6. std::string lastKey; if (const char* c = param("cursor")) lastKey = hexDecode(c); const bool haveCursor = !lastKey.empty(); // Диспетчеризация режимов по приоритету (§3): top > from|to > все. const char* topStr = param("top"); const char* fromStr = param("from"); const char* toStr = param("to"); const char* limitStr = param("limit"); const bool hasTop = (topStr != nullptr); const bool hasFrom = (fromStr != nullptr); const bool hasTo = (toStr != nullptr); const bool hasRange = hasFrom || hasTo; // pageSize (§4.1). Для top параметр ?limit= не применяется. size_t pageSize = kHttpPageSize; if (!hasTop && limitStr) { size_t l = parseUIntParam(limitStr); if (l > 0) pageSize = std::min(l, kHttpPageSizeMax); } std::vector<std::pair<std::string, std::string>> items; std::string nextCursor; if (hasTop) { const size_t N = parseUIntParam(topStr); if (N > 0) { if (haveCursor) { // Сканируем с начала, пропуская ключи <= lastKey (счёт = alreadyReturned), // затем собираем min(pageSize, N-alreadyReturned)+1 для детекции «есть ещё». size_t skipped = 0; size_t alreadyReturned = 0; size_t pageSizeThisPage = std::min(pageSize, N); bool pastCursor = false; bool nReached = false; ScanCallback cb = [&](const std::string& key, const std::string& value) -> bool { if (!pastCursor) { if (key <= lastKey) { ++skipped; return true; } pastCursor = true; alreadyReturned = skipped; if (alreadyReturned >= N) { nReached = true; return false; } pageSizeThisPage = std::min(pageSize, N - alreadyReturned); } items.emplace_back(key, value); if (items.size() >= pageSizeThisPage + 1) return false; return true; }; db_.scanAll(0, cb); if (!nReached && items.size() > pageSizeThisPage) { items.resize(pageSizeThisPage); if (alreadyReturned + pageSizeThisPage < N) { nextCursor = hexEncode(items.back().first); } } } else { size_t pageSizeThisPage = std::min(pageSize, N); ScanCallback cb = [&](const std::string& key, const std::string& value) -> bool { items.emplace_back(key, value); if (items.size() >= pageSizeThisPage + 1) return false; return true; }; db_.scanAll(0, cb); if (items.size() > pageSizeThisPage) { items.resize(pageSizeThisPage); if (pageSizeThisPage < N) { nextCursor = hexEncode(items.back().first); } } } } } else if (hasRange) { std::string from = fromStr ? std::string(fromStr) : ""; std::string to = toStr ? std::string(toStr) : ""; const size_t pageSizeThisPage = pageSize; if (!to.empty()) { ScanCallback cb = [&](const std::string& key, const std::string& value) -> bool { if (haveCursor && key <= lastKey) return true; items.emplace_back(key, value); if (items.size() >= pageSizeThisPage + 1) return false; return true; }; db_.scanStream(from, to, 0, cb); } else { ScanCallback cb = [&](const std::string& key, const std::string& value) -> bool { if (key < from) return true; if (haveCursor && key <= lastKey) return true; items.emplace_back(key, value); if (items.size() >= pageSizeThisPage + 1) return false; return true; }; db_.scanAll(0, cb); } if (items.size() > pageSizeThisPage) { items.resize(pageSizeThisPage); nextCursor = hexEncode(items.back().first); } } else { // Режим «все»: выбирается только отсутствием top/from/to (§3.3). const size_t pageSizeThisPage = pageSize; ScanCallback cb = [&](const std::string& key, const std::string& value) -> bool { if (haveCursor && key <= lastKey) return true; items.emplace_back(key, value); if (items.size() >= pageSizeThisPage + 1) return false; return true; }; db_.scanAll(0, cb); if (items.size() > pageSizeThisPage) { items.resize(pageSizeThisPage); nextCursor = hexEncode(items.back().first); } } // Сериализация ответа (§5). std::ostringstream oss; oss << "{\"items\":["; for (size_t i = 0; i < items.size(); ++i) { if (i) oss << ","; oss << "{\"key\":\"" << jsonEscape(items[i].first) << "\",\"value\":\"" << jsonEscape(items[i].second) << "\"}"; } oss << "],\"count\":" << items.size() << ",\"nextCursor\":\"" << nextCursor << "\"}"; return sendJson(connection, oss.str(), 200); } //============================================================================ // PUT /kv/{key} body=value //============================================================================ MHD_Result HttpHandler::processPut(MHD_Connection* connection, const std::string& url, const std::string& body) { if (url.compare(0, 4, "/kv/") != 0) { return sendError(connection, "Not found", 404); } std::string key = urlDecode(url.substr(4)); if (db_.put(key, body)) { return sendStatus(connection, 204); } return sendError(connection, "Put failed (invalid key/value size)", 400); } //============================================================================ // DELETE /kv/{key} //============================================================================ MHD_Result HttpHandler::processDelete(MHD_Connection* connection, const std::string& url) { if (url.compare(0, 4, "/kv/") != 0) { return sendError(connection, "Not found", 404); } std::string key = urlDecode(url.substr(4)); if (db_.remove(key)) { return sendStatus(connection, 204); } return sendError(connection, "Delete failed", 400); } //============================================================================ // POST /admin/flush | /admin/compact //============================================================================ MHD_Result HttpHandler::processPost(MHD_Connection* connection, const std::string& url) { std::string path = url; auto q = path.find('?'); if (q != std::string::npos) path = path.substr(0, q); if (path == "/admin/flush") { return db_.flush() ? sendStatus(connection, 204) : sendError(connection, "flush failed", 500); } if (path == "/admin/compact") { return db_.compact() ? sendStatus(connection, 204) : sendError(connection, "compact failed", 500); } return sendError(connection, "Not found", 404); } //============================================================================ // Вспомогательные отправители //============================================================================ MHD_Result HttpHandler::sendStatus(MHD_Connection* connection, int status) { MHD_Response* response = MHD_create_response_from_buffer( 0, nullptr, MHD_RESPMEM_PERSISTENT); MHD_Result ret = MHD_queue_response(connection, status, response); MHD_destroy_response(response); return ret; } MHD_Result HttpHandler::sendJson(MHD_Connection* connection, const std::string& json, int status) { MHD_Response* response = MHD_create_response_from_buffer( json.size(), const_cast<char*>(json.data()), MHD_RESPMEM_MUST_COPY); MHD_add_response_header(response, "Content-Type", "application/json"); MHD_add_response_header(response, "Access-Control-Allow-Origin", "*"); MHD_Result ret = MHD_queue_response(connection, status, response); MHD_destroy_response(response); return ret; } MHD_Result HttpHandler::sendBinary(MHD_Connection* connection, const std::string& body, int status) { MHD_Response* response = MHD_create_response_from_buffer( body.size(), const_cast<char*>(body.data()), MHD_RESPMEM_MUST_COPY); MHD_add_response_header(response, "Content-Type", "application/octet-stream"); MHD_Result ret = MHD_queue_response(connection, status, response); MHD_destroy_response(response); return ret; } MHD_Result HttpHandler::sendError(MHD_Connection* connection, const std::string& message, int status) { std::string json = "{\"error\":\"" + jsonEscape(message) + "\"}"; return sendJson(connection, json, status); } }