]> git-server-git.apps.pok.os.sepia.ceph.com Git - ceph.git/commitdiff
nvmeof: fix data races on monitor-client state shared with the timer 69785/head
authorKefu Chai <k.chai@proxmox.com>
Sun, 28 Jun 2026 09:07:22 +0000 (17:07 +0800)
committerKefu Chai <k.chai@proxmox.com>
Mon, 6 Jul 2026 12:31:54 +0000 (20:31 +0800)
send_beacon() runs on the timer thread under beacon_lock, while
handle_nvmeof_gw_map() runs on the dispatch thread under lock. both
touch osdmap_epoch, gwmap_epoch, last_map_time, set_group_id,
cluster_beacon_diff_included and the gw map, each written under one lock
and read under the other.

only the gw map race matters in practice. send_beacon walks it with
get_gw_state() while handle reassigns it with map = new_map, so a beacon
can follow a freed node and crash the ceph-nvmeof-monitor-client. the
beacons then stop, and if the restart outlasts the monitor's beacon
grace the gateway is marked down and its namespaces fail over to a
standby. the scalar races are word-sized stale reads with no observed
symptom. nvmeof gateway deployments only.

user i/o is unaffected unless that failover fires, since the block data
path never runs the racy code. on failover multipath hosts briefly
repath to the standby and single-path ones stall on the owned
namespaces, with no data loss or corruption either way.

make the shared scalars atomic and keep the map private to the dispatch
thread. send_beacon only needs to know whether this gateway is in the
map, so it reads an atomic gw_in_map flag that handle publishes instead
of the map.

send_beacon also read cluster_beacon_diff_included twice per beacon, to
pick the payload and to set its format flag, so a concurrent flip could
tag a diff payload as full. snapshot it once.

Signed-off-by: Kefu Chai <k.chai@proxmox.com>
src/nvmeof/NVMeofGwMonitorClient.cc
src/nvmeof/NVMeofGwMonitorClient.h

index 7f2848027823c85fa9e803b6bb8ff87d6f3e8277..acdf1f732be43ea028728542a4fc9b5f17fb1c42 100644 (file)
@@ -45,8 +45,8 @@ NVMeofGwMonitorClient::NVMeofGwMonitorClient(int argc, const char **argv) :
   gwmap_epoch(0),
   last_map_time(std::chrono::steady_clock::now()),
   reset_timestamp(std::chrono::steady_clock::now()),
-  start_time(last_map_time),
-  cluster_beacon_diff_included(0),
+  start_time(last_map_time.load()),
+  cluster_beacon_diff_included(false),
   poolctx(),
   monc{g_ceph_context, poolctx},
   client_messenger(Messenger::create(g_ceph_context, "async", entity_name_t::CLIENT(-1), "client", getpid())),
@@ -249,31 +249,33 @@ void NVMeofGwMonitorClient::send_beacon()
   BeaconSubsystems subsystem_diff = current_subsystems;
   determine_subsystem_changes(prev_beacon_subsystems, subsystem_diff);
 
-  auto group_key = std::make_pair(pool, group);
-  NvmeGwClientState old_gw_state;
-  // if already got gateway state in the map
-  if (first_beacon == false && get_gw_state("old map", map, group_key, name, old_gw_state))
+  // gw_in_map is published by the dispatch thread, so we avoid reading `map` here
+  if (!first_beacon && gw_in_map.load())
     gw_availability = ok ? gw_availability_t::GW_AVAILABLE : gw_availability_t::GW_UNAVAILABLE;
   dout(1) << "sending beacon as gid " << monc.get_global_id() << " availability " << (int)gw_availability <<
-    " osdmap_epoch " << osdmap_epoch << " gwmap_epoch " << gwmap_epoch << dendl;
+    " osdmap_epoch " << osdmap_epoch.load() << " gwmap_epoch " << gwmap_epoch.load() << dendl;
+
+  // snapshot the cluster feature once so the chosen payload and the format
+  // flag below stay consistent if the dispatch thread updates it mid-beacon.
+  const bool diff_included = cluster_beacon_diff_included.load();
 
   // Check if NVMEOF_BEACON_DIFF feature is supported by the cluster
-  dout(10) << fmt::format("NVMEOF_BEACON_DIFF supported: {}",  cluster_beacon_diff_included ? "yes" : "no") << dendl;
+  dout(10) << fmt::format("NVMEOF_BEACON_DIFF supported: {}",  diff_included ? "yes" : "no") << dendl;
 
   // Send beacon with appropriate version based on cluster features
   auto m = ceph::make_message<MNVMeofGwBeacon>(
       name,
       pool,
       group,
-      cluster_beacon_diff_included ? subsystem_diff : current_subsystems,
+      diff_included ? subsystem_diff : current_subsystems,
       gw_availability,
-      osdmap_epoch,
-      gwmap_epoch,
+      osdmap_epoch.load(),
+      gwmap_epoch.load(),
       beacon_sequence,
       // Pass affected features to the constructor
-      cluster_beacon_diff_included ? 1 : 0
+      diff_included ? 1 : 0
   );
-  dout(10) << "sending beacon with diff support: " << (cluster_beacon_diff_included ? "enabled" : "disabled") << dendl;
+  dout(10) << "sending beacon with diff support: " << (diff_included ? "enabled" : "disabled") << dendl;
   
   monc.send_mon_message(std::move(m));
   ++beacon_sequence;
@@ -284,7 +286,7 @@ void NVMeofGwMonitorClient::disconnect_panic()
 {
   auto disconnect_panic_duration = g_conf().get_val<std::chrono::seconds>("nvmeof_mon_client_disconnect_panic").count();
   auto now = std::chrono::steady_clock::now();
-  auto elapsed_seconds = std::chrono::duration_cast<std::chrono::seconds>(now - last_map_time).count();
+  auto elapsed_seconds = std::chrono::duration_cast<std::chrono::seconds>(now - last_map_time.load()).count();
   if (elapsed_seconds > disconnect_panic_duration) {
     dout(4) << "Triggering a panic upon disconnection from the monitor, elapsed " << elapsed_seconds << ", configured disconnect panic duration " << disconnect_panic_duration << dendl;
     throw std::runtime_error("Lost connection to the monitor (beacon timeout).");
@@ -294,7 +296,7 @@ void NVMeofGwMonitorClient::disconnect_panic()
 void NVMeofGwMonitorClient::connect_panic()
 {
   // Return immediately if the gateway was assigned group ID by the monitor
-  if (set_group_id) {
+  if (set_group_id.load()) {
     return;
   }
   // If the gateway has not been assigned a group ID, panic after timeout
@@ -415,10 +417,11 @@ void NVMeofGwMonitorClient::handle_nvmeof_gw_map(ceph::ref_t<MNVMeofGwMap> nmap)
   // ensure that the gateway state has not vanished
   ceph_assert(got_new_gw_state || !got_old_gw_state);
 
-  uint64_t old_cluster_beacon_diff_included = cluster_beacon_diff_included;
+  bool old_cluster_beacon_diff_included = cluster_beacon_diff_included.load();
   cluster_beacon_diff_included = (new_gw_state.map_features & NVMeofGwMap::FLAG_BEACONDIFF) != 0;
-  if (old_cluster_beacon_diff_included != cluster_beacon_diff_included) {
-    dout(0) << fmt::format("Updated cluster features: 0x{:x}", cluster_beacon_diff_included)
+  if (old_cluster_beacon_diff_included != cluster_beacon_diff_included.load()) {
+    dout(0) << fmt::format("Updated cluster features: NVMEOF_BEACON_DIFF {}",
+                           cluster_beacon_diff_included.load() ? "enabled" : "disabled")
             << dendl;
   }
 
@@ -443,13 +446,13 @@ void NVMeofGwMonitorClient::handle_nvmeof_gw_map(ceph::ref_t<MNVMeofGwMap> nmap)
       dout(10) << "Can not find new gw state" << dendl;
       return;
     }
-    ceph_assert(!set_group_id);
-    while (!set_group_id) {
+    ceph_assert(!set_group_id.load());
+    while (!set_group_id.load()) {
       NVMeofGwMonitorGroupClient monitor_group_client(
           grpc::CreateChannel(monitor_address, gw_creds()));
       dout(10) << "GRPC set_group_id: " <<  new_gw_state.group_id << dendl;
       set_group_id = monitor_group_client.set_group_id( new_gw_state.group_id);
-      if (!set_group_id) {
+      if (!set_group_id.load()) {
              dout(10) << "GRPC set_group_id failed" << dendl;
              auto retry_timeout = g_conf().get_val<uint64_t>("mon_nvmeofgw_set_group_id_retry");
              usleep(retry_timeout);
@@ -548,12 +551,15 @@ void NVMeofGwMonitorClient::handle_nvmeof_gw_map(ceph::ref_t<MNVMeofGwMap> nmap)
       }
     }
     // Update latest accepted osdmap epoch, for beacons
-    if (max_blocklist_epoch > osdmap_epoch) {
+    if (max_blocklist_epoch > osdmap_epoch.load()) {
       osdmap_epoch = max_blocklist_epoch;
-      dout(10) << "Ready for blocklist osd map epoch: " << osdmap_epoch << dendl;
+      dout(10) << "Ready for blocklist osd map epoch: " << osdmap_epoch.load() << dendl;
     }
   }
   map = new_map;
+  // publish, for the timer thread, whether this gateway is present in the
+  // committed map; see gw_in_map and send_beacon().
+  gw_in_map = got_new_gw_state;
 }
 
 Dispatcher::dispatch_result_t NVMeofGwMonitorClient::ms_dispatch2(const ref_t<Message>& m)
index 217a70441470edd3501e39aa6caa4cba8644f7e9..8d9125ab2762b7bd85c1254a7c41f55cdd3cfc4f 100644 (file)
@@ -30,6 +30,8 @@
 #include <grpcpp/grpcpp.h>
 #include <grpcpp/security/credentials.h>
 
+#include <atomic>
+
 class NVMeofGwMonitorClient: public Dispatcher,
                   public md_config_obs_t {
 private:
@@ -43,9 +45,11 @@ private:
   std::string client_cert;
   grpc::SslCredentialsOptions
               gw_ssl_opts;  // gateway grpc ssl options
-  epoch_t     osdmap_epoch; // last awaited osdmap_epoch
-  epoch_t     gwmap_epoch;  // last received gw map epoch
-  std::chrono::time_point<std::chrono::steady_clock>
+  // shared with the timer thread (send_beacon/disconnect_panic); atomic so the
+  // two threads need no common lock and `map` stays private to the dispatch thread
+  std::atomic<epoch_t> osdmap_epoch; // last awaited osdmap_epoch
+  std::atomic<epoch_t> gwmap_epoch;  // last received gw map epoch
+  std::atomic<std::chrono::time_point<std::chrono::steady_clock>>
               last_map_time; // used to panic on disconnect
   std::chrono::time_point<std::chrono::steady_clock>
                 reset_timestamp; // used to bypass some validations
@@ -53,10 +57,16 @@ private:
                 start_time; // used to panic on connect
 
   bool first_beacon = true;
-  bool set_group_id = false;
+  // written by the dispatch thread, read by the timer thread (connect_panic)
+  std::atomic<bool> set_group_id = false;
+  // published by the dispatch thread once this gateway appears in the gw map,
+  // read by the timer thread (send_beacon) instead of reading `map` directly
+  std::atomic<bool> gw_in_map = false;
   uint64_t beacon_sequence = 0;
   BeaconSubsystems prev_beacon_subsystems;
-  bool cluster_beacon_diff_included = 0;  // track cluster features for beacon encoding
+  // written by the dispatch thread, read by the timer thread (send_beacon);
+  // tracks cluster features for beacon encoding
+  std::atomic<bool> cluster_beacon_diff_included = false;
   // init gw ssl opts
   void init_gw_ssl_opts();