diff --git a/src/client/CMakeLists.txt b/src/client/CMakeLists.txt index 4354557c..cd339cac 100644 --- a/src/client/CMakeLists.txt +++ b/src/client/CMakeLists.txt @@ -12,7 +12,7 @@ if (RDMACM_LIBRARIES) set(MSGR_RDMACM "msgr_rdmacm.cpp") endif (RDMACM_LIBRARIES) add_library(vitastor_common STATIC - ../util/epoll_manager.cpp etcd_state_client.cpp messenger.cpp msgr_iothread.cpp ../util/addr_util.cpp + ../util/epoll_manager.cpp etcd_state_client.cpp etcd_state_client_http.cpp messenger.cpp msgr_iothread.cpp ../util/addr_util.cpp msgr_stop.cpp msgr_op.cpp msgr_send.cpp msgr_receive.cpp ../util/ringloop.cpp ../../json11/json11.cpp http_client.cpp osd_ops.cpp pg_states.cpp ../util/timerfd_manager.cpp ../util/str_util.cpp ../util/json_util.cpp ${MSGR_RDMA} ${MSGR_RDMACM} ) @@ -22,6 +22,7 @@ target_compile_options(vitastor_common PUBLIC -fPIC) # libvitastor_client.so add_library(vitastor_client SHARED cluster_client.cpp + cluster_client_real.cpp cluster_client_list.cpp cluster_client_wb.cpp vitastor_c.cpp @@ -96,10 +97,9 @@ add_executable(test_cluster_client EXCLUDE_FROM_ALL ../test/test_cluster_client.cpp pg_states.cpp osd_ops.cpp cluster_client.cpp cluster_client_list.cpp cluster_client_wb.cpp msgr_op.cpp ../test/mock/messenger.cpp msgr_stop.cpp - etcd_state_client.cpp ../util/timerfd_manager.cpp ../util/addr_util.cpp ../util/str_util.cpp ../util/json_util.cpp ../../json11/json11.cpp + etcd_state_client.cpp etcd_state_client_mock.cpp ../util/timerfd_manager.cpp ../util/addr_util.cpp ../util/str_util.cpp ../util/json_util.cpp ../../json11/json11.cpp ) target_link_libraries(test_cluster_client ${LIBURING_LIBRARIES}) -target_compile_definitions(test_cluster_client PUBLIC -D__MOCK__) target_include_directories(test_cluster_client BEFORE PUBLIC ${CMAKE_SOURCE_DIR}/src/test/mock) add_dependencies(build_tests test_cluster_client) add_test(NAME test_cluster_client COMMAND test_cluster_client) diff --git a/src/client/cluster_client.cpp b/src/client/cluster_client.cpp index 44f845e4..c4b1ef13 100644 --- a/src/client/cluster_client.cpp +++ b/src/client/cluster_client.cpp @@ -11,7 +11,7 @@ #define TRY_SEND_CONNECTING 1 #define TRY_SEND_OK 2 -cluster_client_t::cluster_client_t(ring_loop_t *ringloop, timerfd_manager_t *tfd, json11::Json config) +cluster_client_t::cluster_client_t(ring_loop_t *ringloop, timerfd_manager_t *tfd, json11::Json config, std::unique_ptr st_cli_ptr) { wb = new writeback_cache_t(); @@ -53,8 +53,7 @@ cluster_client_t::cluster_client_t(ring_loop_t *ringloop, timerfd_manager_t *tfd }; msgr.parse_config(config); - st_cli = std::make_unique(); - st_cli->tfd = tfd; + st_cli = std::move(st_cli_ptr); st_cli->on_load_config_hook = [this](json11::Json::object & cfg) { on_load_config_hook(cfg); }; st_cli->on_change_osd_state_hook = [this](uint64_t peer_osd) { on_change_osd_state_hook(peer_osd); }; st_cli->on_change_pool_config_hook = [this]() { on_change_pool_config_hook(); }; @@ -62,7 +61,7 @@ cluster_client_t::cluster_client_t(ring_loop_t *ringloop, timerfd_manager_t *tfd st_cli->on_change_pg_state_hook = [this](pool_id_t pool_id, pg_num_t pg_num, osd_num_t prev_primary) { on_change_pg_state_hook(pool_id, pg_num, prev_primary); }; st_cli->on_change_node_placement_hook = [this]() { on_change_node_placement_hook(); }; st_cli->on_load_pgs_hook = [this](bool success) { on_load_pgs_hook(success); }; - st_cli->on_reload_hook = [this]() { st_cli->load_global_config(); }; + st_cli->on_reload_hook = [this]() { this->st_cli->load_global_config(); }; st_cli->parse_config(config); st_cli->infinite_start = false; diff --git a/src/client/cluster_client.h b/src/client/cluster_client.h index d7819c4d..68b14965 100644 --- a/src/client/cluster_client.h +++ b/src/client/cluster_client.h @@ -4,7 +4,7 @@ #pragma once #include "messenger.h" -#include "etcd_state_client.h" +#include "etcd_state_client_http.h" #define DEFAULT_CLIENT_MAX_DIRTY_BYTES 32*1024*1024 #define DEFAULT_CLIENT_MAX_DIRTY_OPS 1024 @@ -139,7 +139,8 @@ public: json11::Json::object cli_config, file_config, etcd_global_config; json11::Json::object config; - cluster_client_t(ring_loop_t *ringloop, timerfd_manager_t *tfd, json11::Json config); + static cluster_client_t* create(ring_loop_t *ringloop, timerfd_manager_t *tfd, json11::Json config); + cluster_client_t(ring_loop_t *ringloop, timerfd_manager_t *tfd, json11::Json config, std::unique_ptr st_cli); ~cluster_client_t(); void execute(cluster_op_t *op); void execute_raw(osd_num_t osd_num, osd_op_t *op); diff --git a/src/client/cluster_client_real.cpp b/src/client/cluster_client_real.cpp new file mode 100644 index 00000000..528fa5d6 --- /dev/null +++ b/src/client/cluster_client_real.cpp @@ -0,0 +1,11 @@ +// Copyright (c) Vitaliy Filippov, 2019+ +// License: VNPL-1.1 or GNU GPL-2.0+ (see README.md for details) + +#include "cluster_client.h" +#include "etcd_state_client_http.h" + +cluster_client_t* cluster_client_t::create(ring_loop_t *ringloop, timerfd_manager_t *tfd, json11::Json config) +{ + auto st_cli = new etcd_state_client_http_t(tfd); + return new cluster_client_t(ringloop, tfd, config, std::unique_ptr(st_cli)); +} diff --git a/src/client/etcd_state_client.cpp b/src/client/etcd_state_client.cpp index 6e4a7da1..e7f7b0fb 100644 --- a/src/client/etcd_state_client.cpp +++ b/src/client/etcd_state_client.cpp @@ -1,13 +1,12 @@ // Copyright (c) Vitaliy Filippov, 2019+ // License: VNPL-1.1 or GNU GPL-2.0+ (see README.md for details) +#include + #include "osd_ops.h" #include "pg_states.h" #include "etcd_state_client.h" -#ifndef __MOCK__ #include "addr_util.h" -#include "http_client.h" -#endif #include "str_util.h" etcd_state_client_t::~etcd_state_client_t() @@ -17,28 +16,8 @@ etcd_state_client_t::~etcd_state_client_t() delete watch; } watches.clear(); - etcd_watches_initialised = -1; -#ifndef __MOCK__ - stop_ws_keepalive(); - if (etcd_watch_ws) - { - http_close(etcd_watch_ws); - etcd_watch_ws = NULL; - } - if (keepalive_client) - { - http_close(keepalive_client); - keepalive_client = NULL; - } -#endif - if (load_pgs_timer_id >= 0) - { - tfd->clear_timer(load_pgs_timer_id); - load_pgs_timer_id = -1; - } } -#ifndef __MOCK__ etcd_kv_t etcd_state_client_t::parse_etcd_kv(const json11::Json & kv_json) { etcd_kv_t kv; @@ -72,104 +51,6 @@ std::vector etcd_state_client_t::get_addresses() return addrs; } -void etcd_state_client_t::etcd_call_oneshot(std::string etcd_address, std::string api, json11::Json payload, - int timeout, std::function callback) -{ - std::string etcd_api_path; - int pos = etcd_address.find('/'); - if (pos >= 0) - { - etcd_api_path = etcd_address.substr(pos); - etcd_address = etcd_address.substr(0, pos); - } - std::string req = payload.dump(); - req = "POST "+etcd_api_path+api+" HTTP/1.1\r\n" - "Host: "+etcd_address+"\r\n" - "Content-Type: application/json\r\n" - "Content-Length: "+std::to_string(req.size())+"\r\n" - "Connection: close\r\n" - "\r\n"+req; - auto http_cli = http_init(tfd); - auto cb = [http_cli, callback](const http_response_t *response) - { - std::string err; - json11::Json data; - response->parse_json_response(err, data); - callback(err, data); - http_close(http_cli); - }; - http_request(http_cli, etcd_address, req, { .timeout = timeout }, cb); -} - -void etcd_state_client_t::etcd_call(std::string api, json11::Json payload, int timeout, - int retries, int interval, std::function callback) -{ - if (!etcd_addresses.size() && !etcd_local.size()) - { - fprintf(stderr, "etcd_address is missing in Vitastor configuration\n"); - exit(1); - } - pick_next_etcd(); - std::string etcd_address = selected_etcd_address; - std::string etcd_api_path; - int pos = etcd_address.find('/'); - if (pos >= 0) - { - etcd_api_path = etcd_address.substr(pos); - etcd_address = etcd_address.substr(0, pos); - } - std::string req = payload.dump(); - req = "POST "+etcd_api_path+api+" HTTP/1.1\r\n" - "Host: "+etcd_address+"\r\n" - "Content-Type: application/json\r\n" - "Content-Length: "+std::to_string(req.size())+"\r\n" - "Connection: keep-alive\r\n" - "Keep-Alive: timeout="+std::to_string(etcd_keepalive_timeout)+"\r\n" - "\r\n"+req; - retries--; - auto cb = [this, api, payload, timeout, retries, interval, callback, - cur_addr = selected_etcd_address](const http_response_t *response) - { - std::string err; - json11::Json data; - response->parse_json_response(err, data); - if (err != "") - { - if (cur_addr == selected_etcd_address) - selected_etcd_address = ""; - if (retries > 0) - { - if (this->log_level > 0) - { - fprintf( - stderr, "Warning: etcd request failed: %s, retrying %d more times\n", - err.c_str(), retries - ); - } - if (interval > 0) - { - // FIXME: Prevent destruction of etcd_state_client if timers or requests are active - tfd->set_timer(interval, false, [this, api, payload, timeout, retries, interval, callback](int) - { - etcd_call(api, payload, timeout, retries, interval, callback); - }); - } - else - etcd_call(api, payload, timeout, retries, interval, callback); - } - else - callback(err, data); - } - else - callback(err, data); - }; - if (!keepalive_client) - { - keepalive_client = http_init(tfd); - } - http_request(keepalive_client, etcd_address, req, { .timeout = timeout, .keepalive = true }, cb); -} - void etcd_state_client_t::add_etcd_url(std::string addr) { if (addr.length() > 0) @@ -256,7 +137,6 @@ void etcd_state_client_t::parse_config(const json11::Json & config) if (this->etcd_keepalive_timeout < 30) this->etcd_keepalive_timeout = 30; } - auto old_etcd_ws_keepalive_interval = this->etcd_ws_keepalive_interval; this->etcd_ws_keepalive_interval = config["etcd_ws_keepalive_interval"].uint64_value(); if (this->etcd_ws_keepalive_interval <= 0) { @@ -282,294 +162,9 @@ void etcd_state_client_t::parse_config(const json11::Json & config) { this->etcd_min_reload_interval = 50; } - if (this->etcd_ws_keepalive_interval != old_etcd_ws_keepalive_interval && ws_keepalive_timer >= 0) - { -#ifndef __MOCK__ - stop_ws_keepalive(); - start_ws_keepalive(); -#endif - } } -void etcd_state_client_t::pick_next_etcd() -{ - if (selected_etcd_address != "") - return; - if (addresses_to_try.size() == 0) - { - // Prefer local etcd, if any - for (int i = 0; i < etcd_local.size(); i++) - addresses_to_try.push_back(etcd_local[i]); - std::vector ns; - for (int i = 0; i < etcd_addresses.size(); i++) - ns.push_back(i); - if (!rand_initialized) - { - timespec tv; - clock_gettime(CLOCK_REALTIME, &tv); - srand48(tv.tv_sec*1000000000 + tv.tv_nsec); - rand_initialized = true; - } - while (ns.size()) - { - int i = lrand48() % ns.size(); - addresses_to_try.push_back(etcd_addresses[ns[i]]); - ns.erase(ns.begin()+i, ns.begin()+i+1); - } - } - selected_etcd_address = addresses_to_try[0]; - addresses_to_try.erase(addresses_to_try.begin(), addresses_to_try.begin()+1); -} - -void etcd_state_client_t::start_etcd_watcher() -{ - if (!etcd_addresses.size() && !etcd_local.size()) - { - fprintf(stderr, "etcd_address is missing in Vitastor configuration\n"); - exit(1); - } - pick_next_etcd(); - std::string etcd_address = selected_etcd_address; - std::string etcd_api_path; - int pos = etcd_address.find('/'); - if (pos >= 0) - { - etcd_api_path = etcd_address.substr(pos); - etcd_address = etcd_address.substr(0, pos); - } - etcd_watches_initialised = 0; - ws_alive = 1; - if (etcd_watch_ws) - { - http_close(etcd_watch_ws); - etcd_watch_ws = NULL; - } - if (this->log_level > 1) - { - fprintf(stderr, "Trying to connect to etcd websocket at %s, watch from revision %ju/%ju/%ju\n", etcd_address.c_str(), - etcd_watch_revision_config, etcd_watch_revision_osd, etcd_watch_revision_pg); - } - etcd_watch_ws = open_websocket(tfd, etcd_address, etcd_api_path+"/watch", etcd_slow_timeout, - [this, cur_addr = selected_etcd_address](const http_response_t *msg) - { - if (msg->body.length()) - { - ws_alive = 1; - std::string json_err; - json11::Json data = json11::Json::parse(msg->body, json_err); - if (json_err != "") - { - fprintf(stderr, "Bad JSON in etcd event: %s, ignoring event\n", json_err.c_str()); - } - else - { - uint64_t watch_id = data["result"]["watch_id"].uint64_value(); - if (data["result"]["created"].bool_value()) - { - if (watch_id == ETCD_CONFIG_WATCH_ID || - watch_id == ETCD_PG_STATE_WATCH_ID || - watch_id == ETCD_OSD_STATE_WATCH_ID) - { - etcd_watches_initialised++; - } - if (etcd_watches_initialised == ETCD_TOTAL_WATCHES && this->log_level > 0) - { - fprintf(stderr, "Successfully subscribed to etcd at %s, revision %ju/%ju/%ju\n", cur_addr.c_str(), - etcd_watch_revision_config, etcd_watch_revision_osd, etcd_watch_revision_pg); - } - } - if (data["result"]["canceled"].bool_value()) - { - // etcd watch canceled, maybe because the revision was compacted - if (data["result"]["compact_revision"].uint64_value()) - { - // we may miss events if we proceed - // so we should restart from the beginning if we can - if (on_reload_hook != NULL) - { - // check to not trigger on_reload_hook multiple times - if (etcd_watch_ws != NULL) - { - fprintf(stderr, "Revisions before %ju were compacted by etcd, reloading state\n", - data["result"]["compact_revision"].uint64_value()); - http_close(etcd_watch_ws); - etcd_watch_ws = NULL; - etcd_watch_revision_config = etcd_watch_revision_osd = etcd_watch_revision_pg = 0; - on_reload_hook(); - } - return; - } - else - { - fprintf(stderr, "Revisions before %ju were compacted by etcd, exiting\n", - data["result"]["compact_revision"].uint64_value()); - exit(1); - } - } - else - { - fprintf(stderr, "Watch canceled by etcd, reason: %s, exiting\n", data["result"]["cancel_reason"].string_value().c_str()); - exit(1); - } - } - // Save revision only if it's present in the message - because sometimes etcd sends something without a header, like: - // {"error": {"grpc_code": 14, "http_code": 503, "http_status": "Service Unavailable", "message": "error reading from server: EOF"}} - // Also don't save revision from the initial created: true messages because they always contain the latest revision - if (etcd_watches_initialised == ETCD_TOTAL_WATCHES && - !data["result"]["header"]["revision"].is_null() && - !data["result"]["created"].bool_value()) - { - // Restart watchers from the same revision number as in the last received message, - // not from the next one to protect against revision being split into multiple messages, - // even though etcd guarantees not to do that **within a single watcher** without fragment=true: - // https://etcd.io/docs/v3.5/learning/api_guarantees/#watch-apis - // Revision contents are ALWAYS split into separate messages for different watchers though! - // So generally we have to resume each watcher from its own revision... - // Progress messages may have watch_id=-1 if sent on behalf of multiple watchers though. - // And antietcd has an advanced semantic which merges the same revision for all watchers - // into one message and just omits watch_id. - // So we also have to handle the case where watch_id is -1 or not present (0). - auto watch_rev = data["result"]["header"]["revision"].uint64_value(); - if (!watch_id || watch_id == UINT64_MAX) - etcd_watch_revision_config = etcd_watch_revision_osd = etcd_watch_revision_pg = watch_rev; - else if (watch_id == ETCD_CONFIG_WATCH_ID) - etcd_watch_revision_config = watch_rev; - else if (watch_id == ETCD_PG_STATE_WATCH_ID) - etcd_watch_revision_pg = watch_rev; - else if (watch_id == ETCD_OSD_STATE_WATCH_ID) - etcd_watch_revision_osd = watch_rev; - addresses_to_try.clear(); - } - // First gather all changes into a hash to remove multiple overwrites - std::map changes; - for (auto & ev: data["result"]["events"].array_items()) - { - auto kv = parse_etcd_kv(ev["kv"]); - if (kv.key != "") - { - changes[kv.key] = kv; - } - } - for (auto & kv: changes) - { - if (this->log_level > 3) - { - fprintf(stderr, "Incoming event: %s -> %s\n", kv.first.c_str(), kv.second.value.dump().c_str()); - } - parse_state(kv.second); - } - // React to changes - if (on_change_hook != NULL) - { - on_change_hook(changes); - } - } - } - if (msg->eof) - { - fprintf(stderr, "Disconnected from etcd %s\n", cur_addr.c_str()); - if (cur_addr == selected_etcd_address) - selected_etcd_address = ""; - if (etcd_watch_ws) - { - http_close(etcd_watch_ws); - etcd_watch_ws = NULL; - } - if (etcd_watches_initialised == 0) - { - // Connection not established, retry in - tfd->set_timer(etcd_quick_timeout, false, [this](int) - { - start_etcd_watcher(); - }); - } - else if (etcd_watches_initialised > 0) - { - // Connection was live, retry immediately - etcd_watches_initialised = 0; - start_etcd_watcher(); - } - } - }); - http_post_message(etcd_watch_ws, WS_TEXT, json11::Json(json11::Json::object { - { "create_request", json11::Json::object { - { "key", base64_encode(etcd_prefix+"/config/") }, - { "range_end", base64_encode(etcd_prefix+"/config0") }, - { "start_revision", etcd_watch_revision_config }, - { "watch_id", ETCD_CONFIG_WATCH_ID }, - { "progress_notify", true }, - } } - }).dump()); - http_post_message(etcd_watch_ws, WS_TEXT, json11::Json(json11::Json::object { - { "create_request", json11::Json::object { - { "key", base64_encode(etcd_prefix+"/osd/state/") }, - { "range_end", base64_encode(etcd_prefix+"/osd/state0") }, - { "start_revision", etcd_watch_revision_osd }, - { "watch_id", ETCD_OSD_STATE_WATCH_ID }, - { "progress_notify", true }, - } } - }).dump()); - http_post_message(etcd_watch_ws, WS_TEXT, json11::Json(json11::Json::object { - { "create_request", json11::Json::object { - { "key", base64_encode(etcd_prefix+"/pg/") }, - { "range_end", base64_encode(etcd_prefix+"/pg0") }, - { "start_revision", etcd_watch_revision_pg }, - { "watch_id", ETCD_PG_STATE_WATCH_ID }, - { "progress_notify", true }, - } } - }).dump()); - // FIXME: Do not watch /pg/history/ at all in client code (not in OSD) - if (on_start_watcher_hook) - { - on_start_watcher_hook(etcd_watch_ws); - } - start_ws_keepalive(); -} - -void etcd_state_client_t::stop_ws_keepalive() -{ - if (ws_keepalive_timer >= 0) - { - tfd->clear_timer(ws_keepalive_timer); - ws_keepalive_timer = -1; - } -} - -void etcd_state_client_t::start_ws_keepalive() -{ - if (ws_keepalive_timer < 0) - { - ws_keepalive_timer = tfd->set_timer(etcd_ws_keepalive_interval*1000, true, [this](int) - { - if (!etcd_watch_ws || etcd_watches_initialised < ETCD_TOTAL_WATCHES) - { - // Do nothing - } - else if (!ws_alive) - { - if (this->log_level > 0) - { - fprintf(stderr, "Websocket ping failed, disconnecting from etcd %s\n", selected_etcd_address.c_str()); - } - if (etcd_watch_ws) - { - http_close(etcd_watch_ws); - etcd_watch_ws = NULL; - } - start_etcd_watcher(); - } - else - { - ws_alive = 0; - http_post_message(etcd_watch_ws, WS_TEXT, json11::Json(json11::Json::object { - { "progress_request", json11::Json::object { } } - }).dump()); - } - }); - } -} - -void etcd_state_client_t::load_global_config() +void etcd_state_client_t::load_global_config(std::function cb) { json11::Json::object req = { { "success", json11::Json::array { json11::Json::object { @@ -583,22 +178,12 @@ void etcd_state_client_t::load_global_config() } } }, } } }; - etcd_txn(req, etcd_quick_timeout, max_etcd_attempts, 0, [this](std::string err, json11::Json data) + etcd_txn(req, etcd_quick_timeout, max_etcd_attempts, 0, [this, cb](std::string err, json11::Json data) { if (err != "") { fprintf(stderr, "Error reading configuration from etcd: %s\n", err.c_str()); - if (infinite_start) - { - tfd->set_timer(etcd_slow_timeout, false, [this](int timer_id) - { - load_global_config(); - }); - } - else - { - exit(1); - } + cb(err); return; } json11::Json config_kv = data["responses"][0]["response_range"]["kvs"][0]; @@ -629,28 +214,12 @@ void etcd_state_client_t::load_global_config() parse_state(kv); } on_load_config_hook(global_config); + cb(""); }); } -void etcd_state_client_t::load_pgs() +void etcd_state_client_t::load_pgs(std::function cb) { - timespec tv; - clock_gettime(CLOCK_REALTIME, &tv); - uint64_t ms_passed = (tv.tv_sec-etcd_last_reload.tv_sec)*1000 + (tv.tv_nsec-etcd_last_reload.tv_nsec)/1000000; - if (ms_passed < etcd_min_reload_interval) - { - if (load_pgs_timer_id < 0) - { - load_pgs_timer_id = tfd->set_timer(etcd_min_reload_interval+50-ms_passed, false, [this](int) { load_pgs(); }); - } - return; - } - etcd_last_reload = tv; - if (load_pgs_timer_id >= 0) - { - tfd->clear_timer(load_pgs_timer_id); - load_pgs_timer_id = -1; - } json11::Json::array txn = { json11::Json::object { { "request_range", json11::Json::object { @@ -698,16 +267,13 @@ void etcd_state_client_t::load_pgs() { req["compare"] = checks; } - etcd_txn_slow(req, [this](std::string err, json11::Json data) + etcd_txn_slow(req, [this, cb](std::string err, json11::Json data) { if (err != "") { // Retry indefinitely fprintf(stderr, "Error loading PGs from etcd: %s\n", err.c_str()); - tfd->set_timer(etcd_slow_timeout, false, [this](int timer_id) - { - load_pgs(); - }); + cb(err); return; } if (!data["succeeded"].bool_value()) @@ -735,24 +301,9 @@ void etcd_state_client_t::load_pgs() } clean_nonexistent_pgs(); on_load_pgs_hook(true); - start_etcd_watcher(); + cb(""); }); } -#else -void etcd_state_client_t::parse_config(const json11::Json & config) -{ -} - -void etcd_state_client_t::load_global_config() -{ - json11::Json::object global_config; - on_load_config_hook(global_config); -} - -void etcd_state_client_t::load_pgs() -{ -} -#endif void etcd_state_client_t::reset_pg_exists() { diff --git a/src/client/etcd_state_client.h b/src/client/etcd_state_client.h index a39149cc..bca1d605 100644 --- a/src/client/etcd_state_client.h +++ b/src/client/etcd_state_client.h @@ -103,15 +103,13 @@ protected: std::vector local_ips; std::vector etcd_addresses; std::vector etcd_local; - std::string selected_etcd_address; - std::vector addresses_to_try; std::vector watches; + std::set seen_peers; bool new_pg_config = false; - int ws_keepalive_timer = -1; - int ws_alive = 0; - bool rand_initialized = false; + void add_etcd_url(std::string); - void pick_next_etcd(); + void reset_pg_exists(); + void clean_nonexistent_pgs(); public: int etcd_keepalive_timeout = 30; int etcd_ws_keepalive_interval = 5; @@ -123,21 +121,15 @@ public: uint64_t global_block_size = DEFAULT_BLOCK_SIZE; uint32_t global_bitmap_granularity = DEFAULT_BITMAP_GRANULARITY; uint32_t global_immediate_commit = IMMEDIATE_NONE; - std::string etcd_prefix; int log_level = 0; - timerfd_manager_t *tfd = NULL; - http_co_t *etcd_watch_ws = NULL, *keepalive_client = NULL; - int etcd_watches_initialised = 0; uint64_t etcd_watch_revision_config = 0; uint64_t etcd_watch_revision_osd = 0; uint64_t etcd_watch_revision_pg = 0; - timespec etcd_last_reload = {}; - int load_pgs_timer_id = -1; + std::map pool_config; std::map peer_states; - std::set seen_peers; std::map inode_config; std::map inode_by_name; json11::Json node_placement; @@ -160,24 +152,22 @@ public: json11::Json::object serialize_inode_cfg(inode_config_t *cfg); etcd_kv_t parse_etcd_kv(const json11::Json & kv_json); std::vector get_addresses(); - void etcd_call_oneshot(std::string etcd_address, std::string api, json11::Json payload, int timeout, std::function callback); - void etcd_call(std::string api, json11::Json payload, int timeout, int retries, int interval, std::function callback); + virtual void etcd_call_oneshot(std::string etcd_address, std::string api, json11::Json payload, int timeout, std::function callback) = 0; + virtual void etcd_call(std::string api, json11::Json payload, int timeout, int retries, int interval, std::function callback) = 0; void etcd_txn(json11::Json txn, int timeout, int retries, int interval, std::function callback); void etcd_txn_slow(json11::Json txn, std::function callback); - void start_etcd_watcher(); - void stop_ws_keepalive(); - void start_ws_keepalive(); - void load_global_config(); - void load_pgs(); - void reset_pg_exists(); - void clean_nonexistent_pgs(); + virtual void etcd_add_watch(json11::Json watch) = 0; + void load_global_config(std::function cb); + virtual void load_global_config() = 0; + void load_pgs(std::function cb); + virtual void load_pgs() = 0; void parse_state(const etcd_kv_t & kv); - void parse_config(const json11::Json & config); + virtual void parse_config(const json11::Json & config); void insert_inode_config(const inode_config_t & cfg); inode_watch_t* watch_inode(std::string name); void close_watch(inode_watch_t* watch); int address_count(); - ~etcd_state_client_t(); + virtual ~etcd_state_client_t(); static uint32_t parse_immediate_commit(const std::string & immediate_commit_str, uint32_t default_value); static uint32_t parse_scheme(const std::string & scheme_str); diff --git a/src/client/etcd_state_client_http.cpp b/src/client/etcd_state_client_http.cpp new file mode 100644 index 00000000..8884c222 --- /dev/null +++ b/src/client/etcd_state_client_http.cpp @@ -0,0 +1,487 @@ +// Copyright (c) Vitaliy Filippov, 2019+ +// License: VNPL-1.1 or GNU GPL-2.0+ (see README.md for details) + +#include "etcd_state_client_http.h" +#include "addr_util.h" +#include "http_client.h" +#include "str_util.h" + +etcd_state_client_http_t::etcd_state_client_http_t(timerfd_manager_t *tfd) +{ + this->tfd = tfd; +} + +etcd_state_client_http_t::~etcd_state_client_http_t() +{ + stop_ws_keepalive(); + if (etcd_watch_ws) + { + http_close(etcd_watch_ws); + etcd_watch_ws = NULL; + } + if (keepalive_client) + { + http_close(keepalive_client); + keepalive_client = NULL; + } + if (load_pgs_timer_id >= 0) + { + tfd->clear_timer(load_pgs_timer_id); + load_pgs_timer_id = -1; + } + etcd_watches_initialised = -1; +} + +void etcd_state_client_http_t::etcd_add_watch(json11::Json watch) +{ + if (etcd_watch_ws) + { + http_post_message(etcd_watch_ws, WS_TEXT, watch.dump()); + } +} + +void etcd_state_client_http_t::etcd_call_oneshot(std::string etcd_address, std::string api, json11::Json payload, + int timeout, std::function callback) +{ + std::string etcd_api_path; + int pos = etcd_address.find('/'); + if (pos >= 0) + { + etcd_api_path = etcd_address.substr(pos); + etcd_address = etcd_address.substr(0, pos); + } + std::string req = payload.dump(); + req = "POST "+etcd_api_path+api+" HTTP/1.1\r\n" + "Host: "+etcd_address+"\r\n" + "Content-Type: application/json\r\n" + "Content-Length: "+std::to_string(req.size())+"\r\n" + "Connection: close\r\n" + "\r\n"+req; + auto http_cli = http_init(tfd); + auto cb = [http_cli, callback](const http_response_t *response) + { + std::string err; + json11::Json data; + response->parse_json_response(err, data); + callback(err, data); + http_close(http_cli); + }; + http_request(http_cli, etcd_address, req, { .timeout = timeout }, cb); +} + +void etcd_state_client_http_t::etcd_call(std::string api, json11::Json payload, int timeout, + int retries, int interval, std::function callback) +{ + if (!etcd_addresses.size() && !etcd_local.size()) + { + fprintf(stderr, "etcd_address is missing in Vitastor configuration\n"); + exit(1); + } + pick_next_etcd(); + std::string etcd_address = selected_etcd_address; + std::string etcd_api_path; + int pos = etcd_address.find('/'); + if (pos >= 0) + { + etcd_api_path = etcd_address.substr(pos); + etcd_address = etcd_address.substr(0, pos); + } + std::string req = payload.dump(); + req = "POST "+etcd_api_path+api+" HTTP/1.1\r\n" + "Host: "+etcd_address+"\r\n" + "Content-Type: application/json\r\n" + "Content-Length: "+std::to_string(req.size())+"\r\n" + "Connection: keep-alive\r\n" + "Keep-Alive: timeout="+std::to_string(etcd_keepalive_timeout)+"\r\n" + "\r\n"+req; + retries--; + auto cb = [this, api, payload, timeout, retries, interval, callback, + cur_addr = selected_etcd_address](const http_response_t *response) + { + std::string err; + json11::Json data; + response->parse_json_response(err, data); + if (err != "") + { + if (cur_addr == selected_etcd_address) + selected_etcd_address = ""; + if (retries > 0) + { + if (this->log_level > 0) + { + fprintf( + stderr, "Warning: etcd request failed: %s, retrying %d more times\n", + err.c_str(), retries + ); + } + if (interval > 0) + { + // FIXME: Prevent destruction of etcd_state_client if timers or requests are active + tfd->set_timer(interval, false, [this, api, payload, timeout, retries, interval, callback](int) + { + etcd_call(api, payload, timeout, retries, interval, callback); + }); + } + else + etcd_call(api, payload, timeout, retries, interval, callback); + } + else + callback(err, data); + } + else + callback(err, data); + }; + if (!keepalive_client) + { + keepalive_client = http_init(tfd); + } + http_request(keepalive_client, etcd_address, req, { .timeout = timeout, .keepalive = true }, cb); +} + +void etcd_state_client_http_t::parse_config(const json11::Json & config) +{ + auto old_etcd_ws_keepalive_interval = this->etcd_ws_keepalive_interval; + etcd_state_client_t::parse_config(config); + if (this->etcd_ws_keepalive_interval != old_etcd_ws_keepalive_interval && ws_keepalive_timer >= 0) + { + stop_ws_keepalive(); + start_ws_keepalive(); + } +} + +void etcd_state_client_http_t::pick_next_etcd() +{ + if (selected_etcd_address != "") + return; + if (addresses_to_try.size() == 0) + { + // Prefer local etcd, if any + for (int i = 0; i < etcd_local.size(); i++) + addresses_to_try.push_back(etcd_local[i]); + std::vector ns; + for (int i = 0; i < etcd_addresses.size(); i++) + ns.push_back(i); + if (!rand_initialized) + { + timespec tv; + clock_gettime(CLOCK_REALTIME, &tv); + srand48(tv.tv_sec*1000000000 + tv.tv_nsec); + rand_initialized = true; + } + while (ns.size()) + { + int i = lrand48() % ns.size(); + addresses_to_try.push_back(etcd_addresses[ns[i]]); + ns.erase(ns.begin()+i, ns.begin()+i+1); + } + } + selected_etcd_address = addresses_to_try[0]; + addresses_to_try.erase(addresses_to_try.begin(), addresses_to_try.begin()+1); +} + +void etcd_state_client_http_t::start_etcd_watcher() +{ + if (!etcd_addresses.size() && !etcd_local.size()) + { + fprintf(stderr, "etcd_address is missing in Vitastor configuration\n"); + exit(1); + } + pick_next_etcd(); + std::string etcd_address = selected_etcd_address; + std::string etcd_api_path; + int pos = etcd_address.find('/'); + if (pos >= 0) + { + etcd_api_path = etcd_address.substr(pos); + etcd_address = etcd_address.substr(0, pos); + } + etcd_watches_initialised = 0; + ws_alive = 1; + if (etcd_watch_ws) + { + http_close(etcd_watch_ws); + etcd_watch_ws = NULL; + } + if (this->log_level > 1) + { + fprintf(stderr, "Trying to connect to etcd websocket at %s, watch from revision %ju/%ju/%ju\n", etcd_address.c_str(), + etcd_watch_revision_config, etcd_watch_revision_osd, etcd_watch_revision_pg); + } + etcd_watch_ws = open_websocket(tfd, etcd_address, etcd_api_path+"/watch", etcd_slow_timeout, + [this, cur_addr = selected_etcd_address](const http_response_t *msg) + { + if (msg->body.length()) + { + ws_alive = 1; + std::string json_err; + json11::Json data = json11::Json::parse(msg->body, json_err); + if (json_err != "") + { + fprintf(stderr, "Bad JSON in etcd event: %s, ignoring event\n", json_err.c_str()); + } + else + { + uint64_t watch_id = data["result"]["watch_id"].uint64_value(); + if (data["result"]["created"].bool_value()) + { + if (watch_id == ETCD_CONFIG_WATCH_ID || + watch_id == ETCD_PG_STATE_WATCH_ID || + watch_id == ETCD_OSD_STATE_WATCH_ID) + { + etcd_watches_initialised++; + } + if (etcd_watches_initialised == ETCD_TOTAL_WATCHES && this->log_level > 0) + { + fprintf(stderr, "Successfully subscribed to etcd at %s, revision %ju/%ju/%ju\n", cur_addr.c_str(), + etcd_watch_revision_config, etcd_watch_revision_osd, etcd_watch_revision_pg); + } + } + if (data["result"]["canceled"].bool_value()) + { + // etcd watch canceled, maybe because the revision was compacted + if (data["result"]["compact_revision"].uint64_value()) + { + // we may miss events if we proceed + // so we should restart from the beginning if we can + if (on_reload_hook != NULL) + { + // check to not trigger on_reload_hook multiple times + if (etcd_watch_ws != NULL) + { + fprintf(stderr, "Revisions before %ju were compacted by etcd, reloading state\n", + data["result"]["compact_revision"].uint64_value()); + http_close(etcd_watch_ws); + etcd_watch_ws = NULL; + etcd_watch_revision_config = etcd_watch_revision_osd = etcd_watch_revision_pg = 0; + on_reload_hook(); + } + return; + } + else + { + fprintf(stderr, "Revisions before %ju were compacted by etcd, exiting\n", + data["result"]["compact_revision"].uint64_value()); + exit(1); + } + } + else + { + fprintf(stderr, "Watch canceled by etcd, reason: %s, exiting\n", data["result"]["cancel_reason"].string_value().c_str()); + exit(1); + } + } + // Save revision only if it's present in the message - because sometimes etcd sends something without a header, like: + // {"error": {"grpc_code": 14, "http_code": 503, "http_status": "Service Unavailable", "message": "error reading from server: EOF"}} + // Also don't save revision from the initial created: true messages because they always contain the latest revision + if (etcd_watches_initialised == ETCD_TOTAL_WATCHES && + !data["result"]["header"]["revision"].is_null() && + !data["result"]["created"].bool_value()) + { + // Restart watchers from the same revision number as in the last received message, + // not from the next one to protect against revision being split into multiple messages, + // even though etcd guarantees not to do that **within a single watcher** without fragment=true: + // https://etcd.io/docs/v3.5/learning/api_guarantees/#watch-apis + // Revision contents are ALWAYS split into separate messages for different watchers though! + // So generally we have to resume each watcher from its own revision... + // Progress messages may have watch_id=-1 if sent on behalf of multiple watchers though. + // And antietcd has an advanced semantic which merges the same revision for all watchers + // into one message and just omits watch_id. + // So we also have to handle the case where watch_id is -1 or not present (0). + auto watch_rev = data["result"]["header"]["revision"].uint64_value(); + if (!watch_id || watch_id == UINT64_MAX) + etcd_watch_revision_config = etcd_watch_revision_osd = etcd_watch_revision_pg = watch_rev; + else if (watch_id == ETCD_CONFIG_WATCH_ID) + etcd_watch_revision_config = watch_rev; + else if (watch_id == ETCD_PG_STATE_WATCH_ID) + etcd_watch_revision_pg = watch_rev; + else if (watch_id == ETCD_OSD_STATE_WATCH_ID) + etcd_watch_revision_osd = watch_rev; + addresses_to_try.clear(); + } + // First gather all changes into a hash to remove multiple overwrites + std::map changes; + for (auto & ev: data["result"]["events"].array_items()) + { + auto kv = parse_etcd_kv(ev["kv"]); + if (kv.key != "") + { + changes[kv.key] = kv; + } + } + for (auto & kv: changes) + { + if (this->log_level > 3) + { + fprintf(stderr, "Incoming event: %s -> %s\n", kv.first.c_str(), kv.second.value.dump().c_str()); + } + parse_state(kv.second); + } + // React to changes + if (on_change_hook != NULL) + { + on_change_hook(changes); + } + } + } + if (msg->eof) + { + fprintf(stderr, "Disconnected from etcd %s\n", cur_addr.c_str()); + if (cur_addr == selected_etcd_address) + selected_etcd_address = ""; + if (etcd_watch_ws) + { + http_close(etcd_watch_ws); + etcd_watch_ws = NULL; + } + if (etcd_watches_initialised == 0) + { + // Connection not established, retry in + tfd->set_timer(etcd_quick_timeout, false, [this](int) + { + start_etcd_watcher(); + }); + } + else if (etcd_watches_initialised > 0) + { + // Connection was live, retry immediately + etcd_watches_initialised = 0; + start_etcd_watcher(); + } + } + }); + http_post_message(etcd_watch_ws, WS_TEXT, json11::Json(json11::Json::object { + { "create_request", json11::Json::object { + { "key", base64_encode(etcd_prefix+"/config/") }, + { "range_end", base64_encode(etcd_prefix+"/config0") }, + { "start_revision", etcd_watch_revision_config }, + { "watch_id", ETCD_CONFIG_WATCH_ID }, + { "progress_notify", true }, + } } + }).dump()); + http_post_message(etcd_watch_ws, WS_TEXT, json11::Json(json11::Json::object { + { "create_request", json11::Json::object { + { "key", base64_encode(etcd_prefix+"/osd/state/") }, + { "range_end", base64_encode(etcd_prefix+"/osd/state0") }, + { "start_revision", etcd_watch_revision_osd }, + { "watch_id", ETCD_OSD_STATE_WATCH_ID }, + { "progress_notify", true }, + } } + }).dump()); + http_post_message(etcd_watch_ws, WS_TEXT, json11::Json(json11::Json::object { + { "create_request", json11::Json::object { + { "key", base64_encode(etcd_prefix+"/pg/") }, + { "range_end", base64_encode(etcd_prefix+"/pg0") }, + { "start_revision", etcd_watch_revision_pg }, + { "watch_id", ETCD_PG_STATE_WATCH_ID }, + { "progress_notify", true }, + } } + }).dump()); + // FIXME: Do not watch /pg/history/ at all in client code (not in OSD) + if (on_start_watcher_hook) + { + on_start_watcher_hook(etcd_watch_ws); + } + start_ws_keepalive(); +} + +void etcd_state_client_http_t::stop_ws_keepalive() +{ + if (ws_keepalive_timer >= 0) + { + tfd->clear_timer(ws_keepalive_timer); + ws_keepalive_timer = -1; + } +} + +void etcd_state_client_http_t::start_ws_keepalive() +{ + if (ws_keepalive_timer < 0) + { + ws_keepalive_timer = tfd->set_timer(etcd_ws_keepalive_interval*1000, true, [this](int) + { + if (!etcd_watch_ws || etcd_watches_initialised < ETCD_TOTAL_WATCHES) + { + // Do nothing + } + else if (!ws_alive) + { + if (this->log_level > 0) + { + fprintf(stderr, "Websocket ping failed, disconnecting from etcd %s\n", selected_etcd_address.c_str()); + } + if (etcd_watch_ws) + { + http_close(etcd_watch_ws); + etcd_watch_ws = NULL; + } + start_etcd_watcher(); + } + else + { + ws_alive = 0; + http_post_message(etcd_watch_ws, WS_TEXT, json11::Json(json11::Json::object { + { "progress_request", json11::Json::object { } } + }).dump()); + } + }); + } +} + +void etcd_state_client_http_t::load_global_config() +{ + etcd_state_client_t::load_global_config([this](const std::string & err) + { + if (err != "") + { + fprintf(stderr, "Error reading configuration from etcd: %s\n", err.c_str()); + if (infinite_start) + { + tfd->set_timer(etcd_slow_timeout, false, [this](int timer_id) + { + load_global_config(); + }); + } + else + { + exit(1); + } + } + }); +} + +void etcd_state_client_http_t::load_pgs() +{ + timespec tv; + clock_gettime(CLOCK_REALTIME, &tv); + uint64_t ms_passed = (tv.tv_sec-etcd_last_reload.tv_sec)*1000 + (tv.tv_nsec-etcd_last_reload.tv_nsec)/1000000; + if (ms_passed < etcd_min_reload_interval) + { + if (load_pgs_timer_id < 0) + { + load_pgs_timer_id = tfd->set_timer(etcd_min_reload_interval+50-ms_passed, false, [this](int) { load_pgs(); }); + } + return; + } + etcd_last_reload = tv; + if (load_pgs_timer_id >= 0) + { + tfd->clear_timer(load_pgs_timer_id); + load_pgs_timer_id = -1; + } + etcd_state_client_t::load_pgs([this](const std::string & err) + { + if (err != "") + { + // Retry indefinitely + fprintf(stderr, "Error loading PGs from etcd: %s\n", err.c_str()); + tfd->set_timer(etcd_slow_timeout, false, [this](int timer_id) + { + load_pgs(); + }); + } + else + { + start_etcd_watcher(); + } + }); +} diff --git a/src/client/etcd_state_client_http.h b/src/client/etcd_state_client_http.h new file mode 100644 index 00000000..33019499 --- /dev/null +++ b/src/client/etcd_state_client_http.h @@ -0,0 +1,36 @@ +// Copyright (c) Vitaliy Filippov, 2019+ +// License: VNPL-1.1 or GNU GPL-2.0+ (see README.md for details) + +#pragma once + +#include "etcd_state_client.h" + +struct __attribute__((visibility("default"))) etcd_state_client_http_t: public etcd_state_client_t +{ +protected: + timerfd_manager_t *tfd = NULL; + std::string selected_etcd_address; + std::vector addresses_to_try; + int ws_keepalive_timer = -1; + int ws_alive = 0; + bool rand_initialized = false; + int etcd_watches_initialised = 0; + timespec etcd_last_reload = {}; + int load_pgs_timer_id = -1; + http_co_t *keepalive_client = NULL; + + void pick_next_etcd(); + void start_etcd_watcher(); + void stop_ws_keepalive(); + void start_ws_keepalive(); +public: + http_co_t *etcd_watch_ws = NULL; + etcd_state_client_http_t(timerfd_manager_t *tfd); + void etcd_call_oneshot(std::string etcd_address, std::string api, json11::Json payload, int timeout, std::function callback) override; + void etcd_call(std::string api, json11::Json payload, int timeout, int retries, int interval, std::function callback) override; + void etcd_add_watch(json11::Json watch) override; + void load_global_config() override; + void load_pgs() override; + void parse_config(const json11::Json & config) override; + ~etcd_state_client_http_t(); +}; diff --git a/src/client/etcd_state_client_mock.cpp b/src/client/etcd_state_client_mock.cpp new file mode 100644 index 00000000..1d25aced --- /dev/null +++ b/src/client/etcd_state_client_mock.cpp @@ -0,0 +1,169 @@ +// Copyright (c) Vitaliy Filippov, 2019+ +// License: VNPL-1.1 or GNU GPL-2.0+ (see README.md for details) + +#include +#include "etcd_state_client_mock.h" +#include "str_util.h" + +void etcd_state_client_mock_t::etcd_add_watch(json11::Json watch) +{ +} + +void etcd_state_client_mock_t::etcd_call_oneshot(std::string etcd_address, std::string api, json11::Json payload, + int timeout, std::function callback) +{ +} + +void etcd_state_client_mock_t::pause() +{ + paused = true; +} + +void etcd_state_client_mock_t::resume() +{ + paused = false; + auto queue = std::move(this->queue); + for (auto& req: queue) + { + etcd_call(req.api, req.payload, req.timeout, req.retries, req.interval, req.callback); + } +} + +void etcd_state_client_mock_t::set(const std::string& key, json11::Json data, uint64_t mod_revision, uint64_t lease_id) +{ + if (!mod_revision) + mod_revision = ++this->mod_revision; + this->data[key] = (etcd_mock_key_data_t){ .value = data.dump(), .mod_revision = mod_revision, .lease_id = lease_id }; +} + +void etcd_state_client_mock_t::etcd_call(std::string api, json11::Json payload, int timeout, + int retries, int interval, std::function callback) +{ + if (paused) + { + queue.push_back({ api, payload, timeout, retries, interval, callback }); + return; + } + printf("+ etcd: %s %s\n", api.c_str(), payload.dump().c_str()); + if (api == "/kv/txn") + { + bool ok = true; + for (auto& check: payload["compare"].array_items()) + { + auto key = base64_decode(check["key"].string_value()); + etcd_mock_key_data_t *key_data = data.find(key) != data.end() ? &data.at(key) : NULL; + auto target = check["target"].string_value(); + auto res = check["result"].string_value(); + assert(res == "LESS" || res == ""); + bool less = res == "LESS"; + if (target == "MOD") + { + uint64_t rev = check["mod_revision"].uint64_value(); + assert(!less || rev); + ok = ok && (less ? (!key_data || key_data->mod_revision < rev) : (key_data && key_data->mod_revision == rev)); + } + else if (target == "CREATE") + { + uint64_t rev = check["create_revision"].uint64_value(); + assert(rev == 0 && !less); + ok = ok && !key_data; + } + else if (target == "VERSION") + { + uint64_t rev = check["version"].uint64_value(); + assert(rev == 0 && !less); + ok = ok && !key_data; + } + else if (target == "LEASE") + { + assert(!less); + uint64_t lease_id = check["lease"].uint64_value(); + ok = ok && key_data && key_data->lease_id == lease_id; + } + else + assert(0); + } + bool has_mod = false; + for (auto& op: payload[ok ? "success" : "failure"].array_items()) + { + auto& obj = op.object_items(); + has_mod = has_mod || obj.find("request_put") != obj.end() || + obj.find("request_delete_range") != obj.end(); + } + if (has_mod) + { + mod_revision++; + } + json11::Json::array responses; + for (auto& op_ptr: payload[ok ? "success" : "failure"].array_items()) + { + auto& op = op_ptr.object_items(); + if (op.find("request_range") != op.end()) + { + json11::Json::array kvs; + auto req = op.at("request_range"); + auto key = base64_decode(req["key"].string_value()); + auto range_end = base64_decode(req["range_end"].string_value()); + auto begin_it = range_end.empty() ? data.find(key) : data.lower_bound(key); + auto end_it = range_end.empty() ? (begin_it == data.end() ? begin_it : std::next(begin_it)) : data.lower_bound(range_end); + for (auto it = begin_it; it != end_it; it++) + { + printf("\\- get: %s = %s, rev %ju\n", it->first.c_str(), it->second.value.c_str(), it->second.mod_revision); + kvs.push_back(json11::Json::object { + { "key", base64_encode(it->first) }, + { "value", base64_encode(it->second.value) }, + { "mod_revision", it->second.mod_revision }, + }); + } + responses.push_back(json11::Json::object { + { "response_range", json11::Json::object{ { "header", json11::Json::object{ { "revision", mod_revision } } }, { "kvs", kvs } } }, + }); + } + else if (op.find("request_put") != op.end()) + { + auto req = op.at("request_put"); + auto key = base64_decode(req["key"].string_value()); + auto value = base64_decode(req["value"].string_value()); + auto lease_id = req["lease"].uint64_value(); + printf("\\- put: %s = %s, rev %ju, lease %ju\n", key.c_str(), value.c_str(), mod_revision, lease_id); + data[key] = { + .value = value, + .mod_revision = mod_revision, + .lease_id = lease_id, + }; + responses.push_back(json11::Json::object { + { "response_put", json11::Json::object{ { "header", json11::Json::object{ { "revision", mod_revision } } } } }, + }); + } + else if (op.find("request_delete_range") != op.end()) + { + auto req = op.at("request_delete_range"); + auto key = base64_decode(req["key"].string_value()); + auto range_end = base64_decode(req["range_end"].string_value()); + uint64_t n_del = 0; + for (auto it = data.lower_bound(key); it != data.end() && (range_end == "" || it->first < range_end); ) + { + printf("\\- del: %s\n", it->first.c_str()); + n_del++; + data.erase(it++); + } + responses.push_back(json11::Json::object { + { "response_delete_range", json11::Json::object{ { "header", json11::Json::object{ { "revision", mod_revision } } }, { "deleted", n_del } } }, + }); + } + } + callback("", json11::Json::object{ { "header", json11::Json::object{ { "revision", mod_revision } } }, { "succeeded", ok }, { "responses", responses } }); + } + else + callback("Unsupported", json11::Json()); +} + +void etcd_state_client_mock_t::load_global_config() +{ + etcd_state_client_t::load_global_config([this](const std::string & err) {}); +} + +void etcd_state_client_mock_t::load_pgs() +{ + etcd_state_client_t::load_pgs([this](const std::string & err) {}); +} diff --git a/src/client/etcd_state_client_mock.h b/src/client/etcd_state_client_mock.h new file mode 100644 index 00000000..e37b3f7f --- /dev/null +++ b/src/client/etcd_state_client_mock.h @@ -0,0 +1,40 @@ +// Copyright (c) Vitaliy Filippov, 2019+ +// License: VNPL-1.1 or GNU GPL-2.0+ (see README.md for details) + +#pragma once + +#include "etcd_state_client.h" + +struct etcd_mock_key_data_t +{ + std::string value; + uint64_t mod_revision; + uint64_t lease_id; +}; + +struct etcd_mock_request_t +{ + std::string api; + json11::Json payload; + int timeout; + int retries; + int interval; + std::function callback; +}; + +struct etcd_state_client_mock_t: public etcd_state_client_t +{ + uint64_t mod_revision = 0; + bool paused = false; + std::vector queue; +public: + std::map data; + void set(const std::string& key, json11::Json data, uint64_t mod_revision = 0, uint64_t lease_id = 0); + void pause(); + void resume(); + void etcd_call_oneshot(std::string etcd_address, std::string api, json11::Json payload, int timeout, std::function callback) override; + void etcd_call(std::string api, json11::Json payload, int timeout, int retries, int interval, std::function callback) override; + void etcd_add_watch(json11::Json watch) override; + void load_global_config() override; + void load_pgs() override; +}; diff --git a/src/client/nbd_proxy.cpp b/src/client/nbd_proxy.cpp index db5d1ef8..aee3d0ce 100644 --- a/src/client/nbd_proxy.cpp +++ b/src/client/nbd_proxy.cpp @@ -514,7 +514,7 @@ help: // Create client ringloop = new ring_loop_t(RINGLOOP_DEFAULT_SIZE); epmgr = new epoll_manager_t(ringloop); - cli = new cluster_client_t(ringloop, epmgr->tfd, cfg); + cli = cluster_client_t::create(ringloop, epmgr->tfd, cfg); if (!inode) { // Load image metadata diff --git a/src/client/ublk_server.cpp b/src/client/ublk_server.cpp index a13f2200..2bac2a94 100644 --- a/src/client/ublk_server.cpp +++ b/src/client/ublk_server.cpp @@ -241,7 +241,7 @@ help: // Create client epmgr = new epoll_manager_t(ringloop); - cli = new cluster_client_t(ringloop, epmgr->tfd, cfg); + cli = cluster_client_t::create(ringloop, epmgr->tfd, cfg); // cli->config contains merged config if (!cfg["queue_depth"].is_null()) diff --git a/src/client/vitastor_c.cpp b/src/client/vitastor_c.cpp index 480e7798..08eef99f 100644 --- a/src/client/vitastor_c.cpp +++ b/src/client/vitastor_c.cpp @@ -103,7 +103,7 @@ vitastor_c *vitastor_c_create_qemu(QEMUSetFDHandler *aio_set_fd_handler, void *a rdma_device, rdma_port_num, rdma_gid_index, rdma_mtu, log_level ); auto self = vitastor_c_create_qemu_common(aio_set_fd_handler, aio_context); - self->cli = new cluster_client_t(NULL, self->tfd, cfg_json); + self->cli = cluster_client_t::create(NULL, self->tfd, cfg_json); return self; } @@ -126,7 +126,7 @@ vitastor_c *vitastor_c_create_qemu_uring(QEMUSetFDHandler *aio_set_fd_handler, v ); auto self = vitastor_c_create_qemu_common(aio_set_fd_handler, aio_context); self->ringloop = ringloop; - self->cli = new cluster_client_t(self->ringloop, self->tfd, cfg_json); + self->cli = cluster_client_t::create(self->ringloop, self->tfd, cfg_json); ringloop->loop(); return self; } @@ -150,7 +150,7 @@ vitastor_c *vitastor_c_create_uring(const char *config_path, const char *etcd_ho vitastor_c *self = new vitastor_c; self->ringloop = ringloop; self->epmgr = new epoll_manager_t(self->ringloop); - self->cli = new cluster_client_t(self->ringloop, self->epmgr->tfd, cfg_json); + self->cli = cluster_client_t::create(self->ringloop, self->epmgr->tfd, cfg_json); ringloop->loop(); return self; } @@ -191,7 +191,7 @@ vitastor_c *vitastor_c_create_uring_json(const char **options, int options_len) vitastor_c *self = new vitastor_c; self->ringloop = ringloop; self->epmgr = new epoll_manager_t(self->ringloop); - self->cli = new cluster_client_t(self->ringloop, self->epmgr->tfd, cfg_json); + self->cli = cluster_client_t::create(self->ringloop, self->epmgr->tfd, cfg_json); ringloop->loop(); return self; } @@ -206,7 +206,7 @@ vitastor_c *vitastor_c_create_epoll_json(const char **options, int options_len) json11::Json cfg_json(cfg); vitastor_c *self = new vitastor_c; self->epmgr = new epoll_manager_t(NULL); - self->cli = new cluster_client_t(NULL, self->epmgr->tfd, cfg_json); + self->cli = cluster_client_t::create(NULL, self->epmgr->tfd, cfg_json); return self; } diff --git a/src/cmd/cli.cpp b/src/cmd/cli.cpp index ae7977f8..e33b7c1e 100644 --- a/src/cmd/cli.cpp +++ b/src/cmd/cli.cpp @@ -555,7 +555,7 @@ static int run(cli_tool_t *p, json11::Json::object cfg) json11::Json cfg_j = cfg; p->ringloop = new ring_loop_t(RINGLOOP_DEFAULT_SIZE); p->epmgr = new epoll_manager_t(p->ringloop); - p->cli = new cluster_client_t(p->ringloop, p->epmgr->tfd, cfg_j); + p->cli = cluster_client_t::create(p->ringloop, p->epmgr->tfd, cfg_j); p->loop_and_wait(action_cb, [&](const cli_result_t & r) { result = r; diff --git a/src/kv/kv_cli.cpp b/src/kv/kv_cli.cpp index fd1caaed..f5c7c709 100644 --- a/src/kv/kv_cli.cpp +++ b/src/kv/kv_cli.cpp @@ -146,7 +146,7 @@ void kv_cli_t::run() // Create client ringloop = new ring_loop_t(512); epmgr = new epoll_manager_t(ringloop); - cli = new cluster_client_t(ringloop, epmgr->tfd, cfg); + cli = cluster_client_t::create(ringloop, epmgr->tfd, cfg); db = new vitastorkv_dbw_t(cli); // Load image metadata while (!cli->is_ready()) diff --git a/src/kv/kv_stress.cpp b/src/kv/kv_stress.cpp index 666b03ed..f530537e 100644 --- a/src/kv/kv_stress.cpp +++ b/src/kv/kv_stress.cpp @@ -272,7 +272,7 @@ void kv_test_t::run(json11::Json cfg) // Create client ringloop = new ring_loop_t(512); epmgr = new epoll_manager_t(ringloop); - cli = new cluster_client_t(ringloop, epmgr->tfd, cfg); + cli = cluster_client_t::create(ringloop, epmgr->tfd, cfg); db = new vitastorkv_dbw_t(cli); // Load image metadata while (!cli->is_ready()) diff --git a/src/nfs/nfs_proxy.cpp b/src/nfs/nfs_proxy.cpp index f0829eca..ac027adc 100644 --- a/src/nfs/nfs_proxy.cpp +++ b/src/nfs/nfs_proxy.cpp @@ -24,7 +24,6 @@ #include "nfs_kv.h" #include "nfs_block.h" #include "nfs_common.h" -#include "http_client.h" #include "cli.h" #define ETCD_INODE_STATS_WATCH_ID 101 @@ -268,7 +267,7 @@ void nfs_proxy_t::run(json11::Json cfg) // Create client ringloop = new ring_loop_t(RINGLOOP_DEFAULT_SIZE); epmgr = new epoll_manager_t(ringloop); - cli = new cluster_client_t(ringloop, epmgr->tfd, cfg); + cli = cluster_client_t::create(ringloop, epmgr->tfd, cfg); cmd = new cli_tool_t(); cmd->ringloop = ringloop; cmd->epmgr = epmgr; @@ -479,28 +478,25 @@ void nfs_proxy_t::watch_stats() parse_stats(kv); } } - if (cli->st_cli->etcd_watch_ws) - { - auto watch_rev = res["header"]["revision"].uint64_value()+1; - http_post_message(cli->st_cli->etcd_watch_ws, WS_TEXT, json11::Json(json11::Json::object { - { "create_request", json11::Json::object { - { "key", base64_encode(cli->st_cli->etcd_prefix+"/inode/stats/") }, - { "range_end", base64_encode(cli->st_cli->etcd_prefix+"/inode/stats0") }, - { "start_revision", watch_rev }, - { "watch_id", ETCD_INODE_STATS_WATCH_ID }, - { "progress_notify", true }, - } } - }).dump()); - http_post_message(cli->st_cli->etcd_watch_ws, WS_TEXT, json11::Json(json11::Json::object { - { "create_request", json11::Json::object { - { "key", base64_encode(cli->st_cli->etcd_prefix+"/pool/stats/") }, - { "range_end", base64_encode(cli->st_cli->etcd_prefix+"/pool/stats0") }, - { "start_revision", watch_rev }, - { "watch_id", ETCD_POOL_STATS_WATCH_ID }, - { "progress_notify", true }, - } } - }).dump()); - } + auto watch_rev = res["header"]["revision"].uint64_value()+1; + cli->st_cli->etcd_add_watch(json11::Json::object { + { "create_request", json11::Json::object { + { "key", base64_encode(cli->st_cli->etcd_prefix+"/inode/stats/") }, + { "range_end", base64_encode(cli->st_cli->etcd_prefix+"/inode/stats0") }, + { "start_revision", watch_rev }, + { "watch_id", ETCD_INODE_STATS_WATCH_ID }, + { "progress_notify", true }, + } } + }); + cli->st_cli->etcd_add_watch(json11::Json::object { + { "create_request", json11::Json::object { + { "key", base64_encode(cli->st_cli->etcd_prefix+"/pool/stats/") }, + { "range_end", base64_encode(cli->st_cli->etcd_prefix+"/pool/stats0") }, + { "start_revision", watch_rev }, + { "watch_id", ETCD_POOL_STATS_WATCH_ID }, + { "progress_notify", true }, + } } + }); }); }; cli->st_cli->on_change_hook = [this, old_hook = cli->st_cli->on_change_hook](std::map & changes) diff --git a/src/osd/osd.cpp b/src/osd/osd.cpp index 1f072b21..6f1c1b8e 100644 --- a/src/osd/osd.cpp +++ b/src/osd/osd.cpp @@ -15,7 +15,7 @@ #include "str_util.h" #include "json_util.h" -osd_t::osd_t(const json11::Json & config, ring_loop_i *ringloop, timerfd_manager_t *tfd) +osd_t::osd_t(const json11::Json & config, ring_loop_i *ringloop, timerfd_manager_t *tfd, std::unique_ptr st_cli_ptr) { zero_buffer_size = 1<<20; zero_buffer = malloc_or_die(zero_buffer_size); @@ -23,6 +23,7 @@ osd_t::osd_t(const json11::Json & config, ring_loop_i *ringloop, timerfd_manager this->ringloop = ringloop; this->tfd = tfd; + this->st_cli = std::move(st_cli_ptr); this->cli_config = config.object_items(); this->file_config = msgr.read_config(this->cli_config); diff --git a/src/osd/osd.h b/src/osd/osd.h index cf6a67c5..0bb14955 100644 --- a/src/osd/osd.h +++ b/src/osd/osd.h @@ -390,7 +390,7 @@ class osd_t } public: - osd_t(const json11::Json & config, ring_loop_i *ringloop, timerfd_manager_t *tfd); + osd_t(const json11::Json & config, ring_loop_i *ringloop, timerfd_manager_t *tfd, std::unique_ptr st_cli_ptr); ~osd_t(); void force_stop(int exitcode); bool shutdown(); diff --git a/src/osd/osd_cluster.cpp b/src/osd/osd_cluster.cpp index ce4e9531..6a8a457a 100644 --- a/src/osd/osd_cluster.cpp +++ b/src/osd/osd_cluster.cpp @@ -3,7 +3,7 @@ #include "osd.h" #include "str_util.h" -#include "etcd_state_client.h" +#include "etcd_state_client_http.h" #include "http_client.h" #include "osd_rmw.h" #include "addr_util.h" @@ -16,7 +16,6 @@ // Peer connection is lost -> Reload connection data -> Try to reconnect void osd_t::init_cluster() { - st_cli = std::make_unique(); if (!st_cli->address_count()) { init_blockstore(NULL); @@ -65,7 +64,6 @@ void osd_t::init_cluster() } else { - st_cli->tfd = tfd; st_cli->log_level = log_level; st_cli->on_change_osd_state_hook = [this](osd_num_t peer_osd) { on_change_osd_state_hook(peer_osd); }; st_cli->on_change_pool_config_hook = [this]() { on_change_pool_config_hook(); }; diff --git a/src/osd/osd_main.cpp b/src/osd/osd_main.cpp index 826190ee..b4f13754 100644 --- a/src/osd/osd_main.cpp +++ b/src/osd/osd_main.cpp @@ -2,6 +2,7 @@ // License: VNPL-1.1 (see README.md for details) #include "epoll_manager.h" +#include "etcd_state_client_http.h" #include "osd.h" #include @@ -65,7 +66,8 @@ int main(int narg, char *args[]) signal(SIGTERM, handle_sigint); ring_loop_t *ringloop = new ring_loop_t(RINGLOOP_DEFAULT_SIZE); epoll_manager_t *epmgr = new epoll_manager_t(ringloop); - osd = new osd_t(config, ringloop, epmgr->tfd); + auto st_cli = new etcd_state_client_http_t(epmgr->tfd); + osd = new osd_t(config, ringloop, epmgr->tfd, std::unique_ptr(st_cli)); while (1) { ringloop->loop(); diff --git a/src/test/test_cas.cpp b/src/test/test_cas.cpp index 9c1e8d26..ea512e02 100644 --- a/src/test/test_cas.cpp +++ b/src/test/test_cas.cpp @@ -70,7 +70,7 @@ int main(int narg, char *args[]) // Create client auto ringloop = new ring_loop_t(RINGLOOP_DEFAULT_SIZE); auto epmgr = new epoll_manager_t(ringloop); - auto cli = new cluster_client_t(ringloop, epmgr->tfd, cfg); + auto cli = cluster_client_t::create(ringloop, epmgr->tfd, cfg); cli->on_ready([&]() { send_read(cli, inode, [&](int r, uint64_t v) diff --git a/src/test/test_cluster_client.cpp b/src/test/test_cluster_client.cpp index c084669c..77ef832f 100644 --- a/src/test/test_cluster_client.cpp +++ b/src/test/test_cluster_client.cpp @@ -4,6 +4,7 @@ #include #include #include +#include "etcd_state_client_mock.h" #include "cluster_client_impl.h" class cluster_client_test_t @@ -15,44 +16,34 @@ public: } }; -void configure_single_pg_pool(cluster_client_t *cli) +void configure_single_pg_pool(etcd_state_client_mock_t *mock) { - cli->st_cli->parse_state((etcd_kv_t){ - .key = "/config/pools", - .value = json11::Json::object { - { "1", json11::Json::object { - { "name", "hddpool" }, - { "scheme", "replicated" }, - { "pg_size", 2 }, - { "pg_minsize", 1 }, - { "pg_count", 1 }, - { "failure_domain", "osd" }, - } } - }, + mock->set("/vitastor/config/pools", json11::Json::object { + { "1", json11::Json::object { + { "name", "hddpool" }, + { "scheme", "replicated" }, + { "pg_size", 2 }, + { "pg_minsize", 1 }, + { "pg_count", 1 }, + { "failure_domain", "osd" }, + { "immediate_commit", "none" }, + } } }); - cli->st_cli->parse_state((etcd_kv_t){ - .key = "/pg/config", - .value = json11::Json::object { - { "items", json11::Json::object { + mock->set("/vitastor/pg/config", json11::Json::object { + { "items", json11::Json::object { + { "1", json11::Json::object { { "1", json11::Json::object { - { "1", json11::Json::object { - { "osd_set", json11::Json::array { 1, 2 } }, - { "primary", 1 }, - } } + { "osd_set", json11::Json::array { 1, 2 } }, + { "primary", 1 }, } } } } - }, + } } }); - cli->st_cli->parse_state((etcd_kv_t){ - .key = "/pg/state/1/1", - .value = json11::Json::object { - { "peers", json11::Json::array { 1, 2 } }, - { "primary", 1 }, - { "state", json11::Json::array { "active" } }, - }, + mock->set("/vitastor/pg/state/1/1", json11::Json::object { + { "peers", json11::Json::array { 1, 2 } }, + { "primary", 1 }, + { "state", json11::Json::array { "active" } }, }); - cli->st_cli->on_load_pgs_hook(true); - cli->st_cli->on_change_pool_config_hook(); } int *test_write(cluster_client_t *cli, uint64_t offset, uint64_t len, uint8_t c, std::function cb = NULL, bool instant = false) @@ -70,7 +61,7 @@ int *test_write(cluster_client_t *cli, uint64_t offset, uint64_t len, uint8_t c, op->callback = [r, cb](cluster_op_t *op) { if (*r == -1) - printf("Error: Not allowed to complete yet\n"); + printf("Error: Not allowed to complete yet (retval %d)\n", op->retval); assert(*r != -1); *r = op->retval == op->len ? 1 : 0; free(op->iov.buf[0].iov_base); @@ -126,15 +117,13 @@ void check_completed(int *r) void pretend_connected(cluster_client_t *cli, osd_num_t osd_num) { printf("OSD %ju connected\n", osd_num); - int peer_fd = cli->msgr.clients.size() ? std::prev(cli->msgr.clients.end())->first+1 : 10; auto cl = new osd_client_t(); cl->client_id = cli->msgr.next_client_id++; cl->osd_num = osd_num; - cl->peer_fd = peer_fd; + cl->peer_fd = -1; cl->peer_state = PEER_CONNECTED; cli->msgr.osd_peers[osd_num] = cl; cli->msgr.clients[cl->client_id] = cl; - cli->msgr.clients_by_fd[peer_fd] = cl; cli->msgr.wanted_peers.erase(osd_num); cli->msgr.repeer_pgs(osd_num); } @@ -209,10 +198,13 @@ void test1() { json11::Json config; timerfd_manager_t *tfd = new timerfd_manager_t([](int fd, bool wr, std::function callback){}); - cluster_client_t *cli = new cluster_client_t(NULL, tfd, config); + etcd_state_client_mock_t *mock = new etcd_state_client_mock_t(); + mock->pause(); + cluster_client_t *cli = new cluster_client_t(NULL, tfd, config, std::unique_ptr(mock)); int *r1 = test_write(cli, 0, 4096, 0x55); - configure_single_pg_pool(cli); + configure_single_pg_pool(mock); + mock->resume(); pretend_connected(cli, 1); can_complete(r1); check_op_count(cli, 1, 1); @@ -421,9 +413,12 @@ void test_writeback() { "client_max_dirty_ops", 2 }, }; timerfd_manager_t *tfd = new timerfd_manager_t([](int fd, bool wr, std::function callback){}); - cluster_client_t *cli = new cluster_client_t(NULL, tfd, config); + etcd_state_client_mock_t *mock = new etcd_state_client_mock_t(); + mock->pause(); + cluster_client_t *cli = new cluster_client_t(NULL, tfd, config, std::unique_ptr(mock)); - configure_single_pg_pool(cli); + configure_single_pg_pool(mock); + mock->resume(); pretend_connected(cli, 1); // Check that 3 consecutive writes are merged by writeback