diff --git a/src/client/cluster_client.cpp b/src/client/cluster_client.cpp index f04dcc68..4f356d69 100644 --- a/src/client/cluster_client.cpp +++ b/src/client/cluster_client.cpp @@ -34,7 +34,7 @@ cluster_client_t::cluster_client_t(ring_loop_t *ringloop, timerfd_manager_t *tfd { // peer_osd just dropped connection // determine WHICH dirty_buffers are now obsolete and repeat them - if (wb->repeat_ops_for(this, peer_osd) > 0) + if (wb->repeat_ops_for(this, peer_osd, 0, 0) > 0) { continue_ops(); } @@ -52,7 +52,8 @@ cluster_client_t::cluster_client_t(ring_loop_t *ringloop, timerfd_manager_t *tfd st_cli.tfd = tfd; 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_hook = [this](std::map & changes) { on_change_hook(changes); }; + st_cli.on_change_pool_config_hook = [this]() { on_change_pool_config_hook(); }; + 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_load_pgs_hook = [this](bool success) { on_load_pgs_hook(success); }; st_cli.on_reload_hook = [this]() { st_cli.load_global_config(); }; @@ -422,7 +423,7 @@ void cluster_client_t::on_load_pgs_hook(bool success) continue_ops(); } -void cluster_client_t::on_change_hook(std::map & changes) +void cluster_client_t::on_change_pool_config_hook() { for (auto pool_item: st_cli.pool_config) { @@ -445,6 +446,19 @@ void cluster_client_t::on_change_hook(std::map & changes continue_ops(); } +void cluster_client_t::on_change_pg_state_hook(pool_id_t pool_id, pg_num_t pg_num, osd_num_t prev_primary) +{ + auto & pg_cfg = st_cli.pool_config[pool_id].pg_config[pg_num]; + if (pg_cfg.cur_primary != prev_primary) + { + // Repeat this PG operations because an OSD which stopped being primary may not fsync operations + if (wb->repeat_ops_for(this, 0, pool_id, pg_num) > 0) + { + continue_ops(); + } + } +} + bool cluster_client_t::get_immediate_commit(uint64_t inode) { if (enable_writeback) @@ -991,6 +1005,29 @@ void cluster_client_t::slice_rw(cluster_op_t *op) } } +bool cluster_client_t::affects_pg(uint64_t inode, uint64_t offset, uint64_t len, pool_id_t pool_id, pg_num_t pg_num) +{ + if (INODE_POOL(inode) != pool_id) + { + return false; + } + auto & pool_cfg = st_cli.pool_config.at(INODE_POOL(inode)); + uint32_t pg_data_size = (pool_cfg.scheme == POOL_SCHEME_REPLICATED ? 1 : pool_cfg.pg_size-pool_cfg.parity_chunks); + uint64_t pg_block_size = pool_cfg.data_block_size * pg_data_size; + uint64_t first_stripe = (offset / pg_block_size) * pg_block_size; + uint64_t last_stripe = len > 0 ? ((offset + len - 1) / pg_block_size) * pg_block_size : first_stripe; + if ((last_stripe/pool_cfg.pg_stripe_size) - (first_stripe/pool_cfg.pg_stripe_size) + 1 >= pool_cfg.real_pg_count) + { + // All PGs are affected + return true; + } + pg_num_t first_pg_num = (first_stripe/pool_cfg.pg_stripe_size) % pool_cfg.real_pg_count + 1; // like map_to_pg() + pg_num_t last_pg_num = (last_stripe/pool_cfg.pg_stripe_size) % pool_cfg.real_pg_count + 1; // like map_to_pg() + return (first_pg_num <= last_pg_num + ? (pg_num >= first_pg_num && pg_num <= last_pg_num) + : (pg_num >= first_pg_num || pg_num <= last_pg_num)); +} + bool cluster_client_t::affects_osd(uint64_t inode, uint64_t offset, uint64_t len, osd_num_t osd) { auto & pool_cfg = st_cli.pool_config.at(INODE_POOL(inode)); diff --git a/src/client/cluster_client.h b/src/client/cluster_client.h index a424404b..fca9fb32 100644 --- a/src/client/cluster_client.h +++ b/src/client/cluster_client.h @@ -79,6 +79,7 @@ class cluster_client_t ring_loop_t *ringloop; std::map pg_counts; + std::map pg_primary; // client_max_dirty_* is actually "max unsynced", for the case when immediate_commit is off uint64_t client_max_dirty_bytes = 0; uint64_t client_max_dirty_ops = 0; @@ -144,9 +145,11 @@ public: protected: bool affects_osd(uint64_t inode, uint64_t offset, uint64_t len, osd_num_t osd); + bool affects_pg(uint64_t inode, uint64_t offset, uint64_t len, pool_id_t pool_id, pg_num_t pg_num); void on_load_config_hook(json11::Json::object & config); void on_load_pgs_hook(bool success); - void on_change_hook(std::map & changes); + void on_change_pool_config_hook(); + void on_change_pg_state_hook(pool_id_t pool_id, pg_num_t pg_num, osd_num_t prev_primary); void on_change_osd_state_hook(uint64_t peer_osd); void execute_internal(cluster_op_t *op); void unshift_op(cluster_op_t *op); diff --git a/src/client/cluster_client_impl.h b/src/client/cluster_client_impl.h index 2f8954b4..78972c26 100644 --- a/src/client/cluster_client_impl.h +++ b/src/client/cluster_client_impl.h @@ -47,7 +47,7 @@ public: bool is_right_merged(dirty_buf_it_t dirty_it); bool is_merged(const dirty_buf_it_t & dirty_it); void copy_write(cluster_op_t *op, int state); - int repeat_ops_for(cluster_client_t *cli, osd_num_t peer_osd); + int repeat_ops_for(cluster_client_t *cli, osd_num_t peer_osd, pool_id_t pool_id, pg_num_t pg_num); void start_writebacks(cluster_client_t *cli, int count); bool read_from_cache(cluster_op_t *op, uint32_t bitmap_granularity); void flush_buffers(cluster_client_t *cli, dirty_buf_it_t from_it, dirty_buf_it_t to_it); diff --git a/src/client/cluster_client_wb.cpp b/src/client/cluster_client_wb.cpp index 997c67b1..86aea2bf 100644 --- a/src/client/cluster_client_wb.cpp +++ b/src/client/cluster_client_wb.cpp @@ -208,7 +208,7 @@ void writeback_cache_t::copy_write(cluster_op_t *op, int state) } } -int writeback_cache_t::repeat_ops_for(cluster_client_t *cli, osd_num_t peer_osd) +int writeback_cache_t::repeat_ops_for(cluster_client_t *cli, osd_num_t peer_osd, pool_id_t pool_id, pg_num_t pg_num) { int repeated = 0; if (dirty_buffers.size()) @@ -218,8 +218,11 @@ int writeback_cache_t::repeat_ops_for(cluster_client_t *cli, osd_num_t peer_osd) for (auto wr_it = dirty_buffers.begin(), flush_it = wr_it, last_it = wr_it; ; ) { bool end = wr_it == dirty_buffers.end(); - bool flush_this = !end && wr_it->second.state != CACHE_REPEATING && - cli->affects_osd(wr_it->first.inode, wr_it->first.stripe, wr_it->second.len, peer_osd); + bool flush_this = !end && wr_it->second.state != CACHE_REPEATING; + if (peer_osd) + flush_this = flush_this && cli->affects_osd(wr_it->first.inode, wr_it->first.stripe, wr_it->second.len, peer_osd); + if (pool_id && pg_num) + flush_this = flush_this && cli->affects_pg(wr_it->first.inode, wr_it->first.stripe, wr_it->second.len, pool_id, pg_num); if (flush_it != wr_it && (end || !flush_this || wr_it->first.inode != flush_it->first.inode || wr_it->first.stripe != last_it->first.stripe+last_it->second.len)) diff --git a/src/client/etcd_state_client.cpp b/src/client/etcd_state_client.cpp index 503d294b..36445683 100644 --- a/src/client/etcd_state_client.cpp +++ b/src/client/etcd_state_client.cpp @@ -890,6 +890,10 @@ void etcd_state_client_t::parse_state(const etcd_kv_t & kv) } } } + if (on_change_pool_config_hook) + { + on_change_pool_config_hook(); + } } else if (key == etcd_prefix+"/config/pgs") { @@ -1028,13 +1032,19 @@ void etcd_state_client_t::parse_state(const etcd_kv_t & kv) else if (value.is_null()) { auto & pg_cfg = this->pool_config[pool_id].pg_config[pg_num]; + auto prev_primary = pg_cfg.cur_primary; pg_cfg.state_exists = false; pg_cfg.cur_primary = 0; pg_cfg.cur_state = 0; + if (on_change_pg_state_hook) + { + on_change_pg_state_hook(pool_id, pg_num, prev_primary); + } } else { auto & pg_cfg = this->pool_config[pool_id].pg_config[pg_num]; + auto prev_primary = pg_cfg.cur_primary; pg_cfg.state_exists = true; osd_num_t cur_primary = value["primary"].uint64_value(); int state = 0; @@ -1065,6 +1075,10 @@ void etcd_state_client_t::parse_state(const etcd_kv_t & kv) } pg_cfg.cur_primary = cur_primary; pg_cfg.cur_state = state; + if (on_change_pg_state_hook) + { + on_change_pg_state_hook(pool_id, pg_num, prev_primary); + } } } else if (key.substr(0, etcd_prefix.length()+11) == etcd_prefix+"/osd/state/") diff --git a/src/client/etcd_state_client.h b/src/client/etcd_state_client.h index 5fcbc3ec..0282934a 100644 --- a/src/client/etcd_state_client.h +++ b/src/client/etcd_state_client.h @@ -127,6 +127,8 @@ public: std::function on_load_config_hook; std::function load_pgs_checks_hook; std::function on_load_pgs_hook; + std::function on_change_pool_config_hook; + std::function on_change_pg_state_hook; std::function on_change_pg_history_hook; std::function on_change_osd_state_hook; std::function on_reload_hook;