From b263d311ef15ba2842086a28dacca9c370c5150a Mon Sep 17 00:00:00 2001 From: Vitaliy Filippov Date: Sat, 20 Jul 2024 17:48:49 +0300 Subject: [PATCH] Use separate watch revisions for different watchers --- src/client/etcd_state_client.cpp | 57 ++++++++++++++++++++++---------- src/client/etcd_state_client.h | 6 ++-- src/cmd/cli_rm_osd.cpp | 2 +- src/nfs/nfs_proxy.cpp | 40 ++++++++++++---------- src/osd/osd_cluster.cpp | 4 +-- 5 files changed, 68 insertions(+), 41 deletions(-) diff --git a/src/client/etcd_state_client.cpp b/src/client/etcd_state_client.cpp index 76ad6cf2..82b518a9 100644 --- a/src/client/etcd_state_client.cpp +++ b/src/client/etcd_state_client.cpp @@ -333,7 +333,10 @@ void etcd_state_client_t::start_etcd_watcher() etcd_watch_ws = NULL; } if (this->log_level > 1) - fprintf(stderr, "Trying to connect to etcd websocket at %s, watch from revision %ju\n", etcd_address.c_str(), etcd_watch_revision); + { + 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) { @@ -348,15 +351,20 @@ void etcd_state_client_t::start_etcd_watcher() } else { + uint64_t watch_id = data["result"]["watch_id"].uint64_value(); if (data["result"]["created"].bool_value()) { - uint64_t watch_id = data["result"]["watch_id"].uint64_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\n", cur_addr.c_str(), etcd_watch_revision); + { + 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()) { @@ -374,7 +382,7 @@ void etcd_state_client_t::start_etcd_watcher() data["result"]["compact_revision"].uint64_value()); http_close(etcd_watch_ws); etcd_watch_ws = NULL; - etcd_watch_revision = 0; + etcd_watch_revision_config = etcd_watch_revision_osd = etcd_watch_revision_pg = 0; on_reload_hook(); } return; @@ -392,13 +400,29 @@ void etcd_state_client_t::start_etcd_watcher() 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"}} if (etcd_watches_initialised == ETCD_TOTAL_WATCHES && !data["result"]["header"]["revision"].is_null()) { - // Protect against a revision being split into multiple messages and some - // of them being lost. - // Also 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"}} - etcd_watch_revision = data["result"]["header"]["revision"].uint64_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 @@ -456,7 +480,7 @@ void etcd_state_client_t::start_etcd_watcher() { "create_request", json11::Json::object { { "key", base64_encode(etcd_prefix+"/config/") }, { "range_end", base64_encode(etcd_prefix+"/config0") }, - { "start_revision", etcd_watch_revision }, + { "start_revision", etcd_watch_revision_config }, { "watch_id", ETCD_CONFIG_WATCH_ID }, { "progress_notify", true }, } } @@ -465,7 +489,7 @@ void etcd_state_client_t::start_etcd_watcher() { "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 }, + { "start_revision", etcd_watch_revision_osd }, { "watch_id", ETCD_OSD_STATE_WATCH_ID }, { "progress_notify", true }, } } @@ -474,7 +498,7 @@ void etcd_state_client_t::start_etcd_watcher() { "create_request", json11::Json::object { { "key", base64_encode(etcd_prefix+"/pg/") }, { "range_end", base64_encode(etcd_prefix+"/pg0") }, - { "start_revision", etcd_watch_revision }, + { "start_revision", etcd_watch_revision_pg }, { "watch_id", ETCD_PG_STATE_WATCH_ID }, { "progress_notify", true }, } } @@ -636,13 +660,10 @@ void etcd_state_client_t::load_pgs() return; } reset_pg_exists(); - if (!etcd_watch_revision) + etcd_watch_revision_config = etcd_watch_revision_osd = etcd_watch_revision_pg = data["header"]["revision"].uint64_value()+1; + if (this->log_level > 3) { - etcd_watch_revision = data["header"]["revision"].uint64_value()+1; - if (this->log_level > 3) - { - fprintf(stderr, "Loaded revision %ju of PG configuration\n", etcd_watch_revision-1); - } + fprintf(stderr, "Loaded revision %ju of PG configuration\n", etcd_watch_revision_pg-1); } for (auto & res: data["responses"].array_items()) { diff --git a/src/client/etcd_state_client.h b/src/client/etcd_state_client.h index fb21771e..5f42e1cb 100644 --- a/src/client/etcd_state_client.h +++ b/src/client/etcd_state_client.h @@ -94,7 +94,6 @@ protected: std::string selected_etcd_address; std::vector addresses_to_try; std::vector watches; - http_co_t *etcd_watch_ws = NULL, *keepalive_client = NULL; bool new_pg_config = false; int ws_keepalive_timer = -1; int ws_alive = 0; @@ -115,8 +114,11 @@ public: 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 = 0; + uint64_t etcd_watch_revision_config = 0; + uint64_t etcd_watch_revision_osd = 0; + uint64_t etcd_watch_revision_pg = 0; std::map pool_config; std::map peer_states; std::set seen_peers; diff --git a/src/cmd/cli_rm_osd.cpp b/src/cmd/cli_rm_osd.cpp index 4c118a22..a7f0a920 100644 --- a/src/cmd/cli_rm_osd.cpp +++ b/src/cmd/cli_rm_osd.cpp @@ -427,7 +427,7 @@ struct rm_osd_t { "target", "MOD" }, { "key", history_key }, { "result", "LESS" }, - { "mod_revision", parent->cli->st_cli.etcd_watch_revision+1 }, + { "mod_revision", parent->cli->st_cli.etcd_watch_revision_pg+1 }, }); } } diff --git a/src/nfs/nfs_proxy.cpp b/src/nfs/nfs_proxy.cpp index 49c2eebf..30d46071 100644 --- a/src/nfs/nfs_proxy.cpp +++ b/src/nfs/nfs_proxy.cpp @@ -372,24 +372,6 @@ void nfs_proxy_t::watch_stats() assert(cli->st_cli.on_start_watcher_hook == NULL); cli->st_cli.on_start_watcher_hook = [this](http_co_t *etcd_watch_ws) { - http_post_message(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", cli->st_cli.etcd_watch_revision }, - { "watch_id", ETCD_INODE_STATS_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(cli->st_cli.etcd_prefix+"/pool/stats/") }, - { "range_end", base64_encode(cli->st_cli.etcd_prefix+"/pool/stats0") }, - { "start_revision", cli->st_cli.etcd_watch_revision }, - { "watch_id", ETCD_POOL_STATS_WATCH_ID }, - { "progress_notify", true }, - } } - }).dump()); cli->st_cli.etcd_txn_slow(json11::Json::object { { "success", json11::Json::array { json11::Json::object { @@ -415,6 +397,28 @@ 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()); + } }); }; cli->st_cli.on_change_hook = [this, old_hook = cli->st_cli.on_change_hook](std::map & changes) diff --git a/src/osd/osd_cluster.cpp b/src/osd/osd_cluster.cpp index 8272cb47..d24c7139 100644 --- a/src/osd/osd_cluster.cpp +++ b/src/osd/osd_cluster.cpp @@ -905,7 +905,7 @@ void osd_t::report_pg_states() { "target", "MOD" }, { "key", state_key_base64 }, { "result", "LESS" }, - { "mod_revision", st_cli.etcd_watch_revision+1 }, + { "mod_revision", st_cli.etcd_watch_revision_pg+1 }, }); continue; } @@ -976,7 +976,7 @@ void osd_t::report_pg_states() { "target", "MOD" }, { "key", history_key }, { "result", "LESS" }, - { "mod_revision", st_cli.etcd_watch_revision+1 }, + { "mod_revision", st_cli.etcd_watch_revision_pg+1 }, }); success.push_back(json11::Json::object { { "request_put", json11::Json::object {