Change st_cli to pointer

This commit is contained in:
Vitaliy Filippov
2026-06-13 19:56:06 +03:00
parent 80fa3094b3
commit 26fb08d7da
46 changed files with 503 additions and 501 deletions
+2 -2
View File
@@ -116,7 +116,7 @@ void osd_t::init_blockstore(std::function<void()> on_init)
auto bs_cfg = json_to_string_map(this->config);
this->bs = blockstore_i::create(bs_cfg, ringloop, tfd);
// Pre-configure pool PG shards
for (auto & pool_item: st_cli.pool_config)
for (auto & pool_item: st_cli->pool_config)
{
auto st = bs->reshard_start(pool_item.first, pool_item.second.pg_count, pool_item.second.pg_stripe_size, 0);
assert(!st);
@@ -164,7 +164,7 @@ void osd_t::parse_config(bool init)
auto bs_cfg = json_to_string_map(config);
bs->parse_config(bs_cfg);
}
st_cli.parse_config(config);
st_cli->parse_config(config);
msgr.parse_config(config);
if (init)
{
+3 -3
View File
@@ -152,7 +152,7 @@ class osd_t
// cluster state
etcd_state_client_t st_cli;
std::unique_ptr<etcd_state_client_t> st_cli;
osd_messenger_t msgr;
int etcd_failed_attempts = 0;
std::string etcd_lease_id;
@@ -383,8 +383,8 @@ class osd_t
inline pg_num_t map_to_pg(object_id oid)
{
auto pool_it = st_cli.pool_config.find(INODE_POOL(oid.inode));
if (pool_it == st_cli.pool_config.end())
auto pool_it = st_cli->pool_config.find(INODE_POOL(oid.inode));
if (pool_it == st_cli->pool_config.end())
return 1;
return (oid.stripe / pool_it->second.applied_pg_stripe_size) % pool_it->second.applied_pg_count + 1;
}
+71 -70
View File
@@ -16,7 +16,8 @@
// Peer connection is lost -> Reload connection data -> Try to reconnect
void osd_t::init_cluster()
{
if (!st_cli.address_count())
st_cli = std::make_unique<etcd_state_client_t>();
if (!st_cli->address_count())
{
init_blockstore(NULL);
if (run_primary)
@@ -30,7 +31,7 @@ void osd_t::init_cluster()
parse_test_peer(pos < 0 ? peerstr : peerstr.substr(0, pos));
peerstr = pos < 0 ? std::string("") : peerstr.substr(pos+1);
}
if (st_cli.peer_states.size() < 2)
if (st_cli->peer_states.size() < 2)
{
throw std::runtime_error("run_primary requires at least 2 peers");
}
@@ -46,7 +47,7 @@ void osd_t::init_cluster()
.target_set = { 1, 2, 3 },
.cur_set = { 0, 0, 0 },
};
st_cli.pool_config[1] = (pool_config_t){
st_cli->pool_config[1] = (pool_config_t){
.exists = true,
.id = 1,
.name = "testpool",
@@ -64,19 +65,19 @@ 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(); };
st_cli.on_change_backfillfull_hook = [this](pool_id_t pool_id) { on_change_backfillfull_hook(pool_id); };
st_cli.on_change_pg_history_hook = [this](pool_id_t pool_id, pg_num_t pg_num) { on_change_pg_history_hook(pool_id, pg_num); };
st_cli.on_change_hook = [this](std::map<std::string, etcd_kv_t> & changes) { on_change_etcd_state_hook(changes); };
st_cli.on_load_config_hook = [this](json11::Json::object & cfg) { on_load_config_hook(cfg); };
st_cli.load_pgs_checks_hook = [this]() { return on_load_pgs_checks_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->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(); };
st_cli->on_change_backfillfull_hook = [this](pool_id_t pool_id) { on_change_backfillfull_hook(pool_id); };
st_cli->on_change_pg_history_hook = [this](pool_id_t pool_id, pg_num_t pg_num) { on_change_pg_history_hook(pool_id, pg_num); };
st_cli->on_change_hook = [this](std::map<std::string, etcd_kv_t> & changes) { on_change_etcd_state_hook(changes); };
st_cli->on_load_config_hook = [this](json11::Json::object & cfg) { on_load_config_hook(cfg); };
st_cli->load_pgs_checks_hook = [this]() { return on_load_pgs_checks_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(); };
peering_state = OSD_LOADING_PGS;
st_cli.load_global_config();
st_cli->load_global_config();
}
if (run_primary && autosync_interval > 0)
{
@@ -100,17 +101,17 @@ void osd_t::parse_test_peer(std::string peer)
osd_num_t peer_osd = strtoull(osd_num_str.c_str(), NULL, 10);
if (!peer_osd)
throw new std::runtime_error("Could not parse OSD peer osd_num");
else if (st_cli.peer_states.find(peer_osd) != st_cli.peer_states.end())
else if (st_cli->peer_states.find(peer_osd) != st_cli->peer_states.end())
throw std::runtime_error("Same osd number "+std::to_string(peer_osd)+" specified twice in peers");
int port = strtoull(port_str.c_str(), NULL, 10);
if (!port)
throw new std::runtime_error("Could not parse OSD peer port");
st_cli.peer_states[peer_osd] = json11::Json::object {
st_cli->peer_states[peer_osd] = json11::Json::object {
{ "state", "up" },
{ "addresses", json11::Json::array { addr } },
{ "port", port },
};
msgr.connect_peer(peer_osd, st_cli.peer_states[peer_osd]);
msgr.connect_peer(peer_osd, st_cli->peer_states[peer_osd]);
}
bool osd_t::check_peer_config(osd_client_t *cl, json11::Json conf)
@@ -349,19 +350,19 @@ void osd_t::report_statistics()
json11::Json::array txn = {
json11::Json::object {
{ "request_put", json11::Json::object {
{ "key", base64_encode(st_cli.etcd_prefix+"/osd/stats/"+std::to_string(osd_num)) },
{ "key", base64_encode(st_cli->etcd_prefix+"/osd/stats/"+std::to_string(osd_num)) },
{ "value", base64_encode(get_statistics().dump()) },
} },
},
json11::Json::object {
{ "request_put", json11::Json::object {
{ "key", base64_encode(st_cli.etcd_prefix+"/osd/space/"+std::to_string(osd_num)) },
{ "key", base64_encode(st_cli->etcd_prefix+"/osd/space/"+std::to_string(osd_num)) },
{ "value", base64_encode(json11::Json(inode_space).dump()) },
} },
},
json11::Json::object {
{ "request_put", json11::Json::object {
{ "key", base64_encode(st_cli.etcd_prefix+"/osd/inodestats/"+std::to_string(osd_num)) },
{ "key", base64_encode(st_cli->etcd_prefix+"/osd/inodestats/"+std::to_string(osd_num)) },
{ "value", base64_encode(json11::Json(inode_ops).dump()) },
} },
},
@@ -385,19 +386,19 @@ void osd_t::report_statistics()
pg_stats["write_osd_set"] = pg.cur_set;
txn.push_back(json11::Json::object {
{ "request_put", json11::Json::object {
{ "key", base64_encode(st_cli.etcd_prefix+"/pgstats/"+std::to_string(pg.pool_id)+"/"+std::to_string(pg.pg_num)) },
{ "key", base64_encode(st_cli->etcd_prefix+"/pgstats/"+std::to_string(pg.pool_id)+"/"+std::to_string(pg.pg_num)) },
{ "value", base64_encode(json11::Json(pg_stats).dump()) },
} }
});
}
st_cli.etcd_txn_slow(json11::Json::object { { "success", txn } }, [this](std::string err, json11::Json res)
st_cli->etcd_txn_slow(json11::Json::object { { "success", txn } }, [this](std::string err, json11::Json res)
{
etcd_reporting_stats = false;
if (err != "")
{
printf("[OSD %ju] Error reporting state to etcd: %s\n", this->osd_num, err.c_str());
// Retry indefinitely
tfd->set_timer(st_cli.etcd_slow_timeout, false, [this](int timer_id)
tfd->set_timer(st_cli->etcd_slow_timeout, false, [this](int timer_id)
{
report_statistics();
});
@@ -412,13 +413,13 @@ void osd_t::report_statistics()
void osd_t::on_change_osd_state_hook(osd_num_t peer_osd)
{
if (st_cli.peer_states[peer_osd].is_null())
if (st_cli->peer_states[peer_osd].is_null())
{
repeer_pgs(peer_osd);
}
else if (msgr.wanted_peers.find(peer_osd) != msgr.wanted_peers.end())
{
msgr.connect_peer(peer_osd, st_cli.peer_states[peer_osd]);
msgr.connect_peer(peer_osd, st_cli->peer_states[peer_osd]);
}
}
@@ -434,8 +435,8 @@ void osd_t::apply_pg_locks_localize_only()
{
for (auto & pp: pgs)
{
auto pool_it = st_cli.pool_config.find(pp.first.pool_id);
if (pool_it == st_cli.pool_config.end())
auto pool_it = st_cli->pool_config.find(pp.first.pool_id);
if (pool_it == st_cli->pool_config.end())
{
continue;
}
@@ -464,19 +465,19 @@ void osd_t::on_change_backfillfull_hook(pool_id_t pool_id)
void osd_t::on_change_etcd_state_hook(std::map<std::string, etcd_kv_t> & changes)
{
if (changes.find(st_cli.etcd_prefix+"/config/global") != changes.end())
if (changes.find(st_cli->etcd_prefix+"/config/global") != changes.end())
{
etcd_global_config = changes[st_cli.etcd_prefix+"/config/global"].value.object_items();
etcd_global_config = changes[st_cli->etcd_prefix+"/config/global"].value.object_items();
parse_config(false);
}
bool pools = changes.find(st_cli.etcd_prefix+"/config/pools") != changes.end();
bool pools = changes.find(st_cli->etcd_prefix+"/config/pools") != changes.end();
if (pools)
{
apply_no_inode_stats();
}
if (run_primary)
{
bool pgs = changes.find(st_cli.etcd_prefix+"/pg/config") != changes.end();
bool pgs = changes.find(st_cli->etcd_prefix+"/pg/config") != changes.end();
if (pools || pgs)
{
apply_pg_count();
@@ -489,7 +490,7 @@ void osd_t::on_load_config_hook(json11::Json::object & global_config)
{
etcd_global_config = global_config;
parse_config(true);
st_cli.on_load_config_hook = [this](json11::Json::object & cfg) { on_reload_config_hook(cfg); };
st_cli->on_load_config_hook = [this](json11::Json::object & cfg) { on_reload_config_hook(cfg); };
etcd_global_config_loaded = true;
init_blockstore([this]()
{
@@ -511,14 +512,14 @@ void osd_t::acquire_lease()
// Apply no_inode_stats before the first statistics report
apply_no_inode_stats();
// Maximum lease TTL is (report interval) + retries * (timeout + repeat interval)
st_cli.etcd_call("/lease/grant", json11::Json::object {
{ "TTL", etcd_report_interval+(st_cli.max_etcd_attempts*(2*st_cli.etcd_quick_timeout)+999)/1000 }
}, st_cli.etcd_quick_timeout, 0, 0, [this](std::string err, json11::Json data)
st_cli->etcd_call("/lease/grant", json11::Json::object {
{ "TTL", etcd_report_interval+(st_cli->max_etcd_attempts*(2*st_cli->etcd_quick_timeout)+999)/1000 }
}, st_cli->etcd_quick_timeout, 0, 0, [this](std::string err, json11::Json data)
{
if (err != "" || data["ID"].string_value() == "")
{
printf("Error acquiring a lease from etcd: %s, retrying\n", err.c_str());
tfd->set_timer(st_cli.etcd_quick_timeout, false, [this](int timer_id)
tfd->set_timer(st_cli->etcd_quick_timeout, false, [this](int timer_id)
{
acquire_lease();
});
@@ -546,9 +547,9 @@ void osd_t::acquire_lease()
// Do it first to allow "monitors" check it when moving PGs
void osd_t::create_osd_state()
{
std::string state_key = base64_encode(st_cli.etcd_prefix+"/osd/state/"+std::to_string(osd_num));
std::string state_key = base64_encode(st_cli->etcd_prefix+"/osd/state/"+std::to_string(osd_num));
self_state = get_osd_state();
st_cli.etcd_txn(json11::Json::object {
st_cli->etcd_txn(json11::Json::object {
// Check that the state key does not exist
{ "compare", json11::Json::array {
json11::Json::object {
@@ -573,19 +574,19 @@ void osd_t::create_osd_state()
} }
},
} },
}, st_cli.etcd_quick_timeout, 0, 0, [this](std::string err, json11::Json data)
}, st_cli->etcd_quick_timeout, 0, 0, [this](std::string err, json11::Json data)
{
if (err != "")
{
etcd_failed_attempts++;
printf("Error creating OSD state key: %s\n", err.c_str());
if (etcd_failed_attempts > st_cli.max_etcd_attempts)
if (etcd_failed_attempts > st_cli->max_etcd_attempts)
{
// Die
throw std::runtime_error("Cluster connection failed");
}
// Retry
tfd->set_timer(st_cli.etcd_quick_timeout, false, [this](int timer_id)
tfd->set_timer(st_cli->etcd_quick_timeout, false, [this](int timer_id)
{
create_osd_state();
});
@@ -594,7 +595,7 @@ void osd_t::create_osd_state()
if (!data["succeeded"].bool_value())
{
// OSD is already up
auto kv = st_cli.parse_etcd_kv(data["responses"][0]["response_range"]["kvs"][0]);
auto kv = st_cli->parse_etcd_kv(data["responses"][0]["response_range"]["kvs"][0]);
printf("Key %s already exists in etcd, OSD %ju is still up\n", kv.key.c_str(), this->osd_num);
int64_t port = kv.value["port"].int64_value();
for (auto & addr: kv.value["addresses"].array_items())
@@ -606,7 +607,7 @@ void osd_t::create_osd_state()
}
if (run_primary)
{
st_cli.load_pgs();
st_cli->load_pgs();
}
report_statistics();
});
@@ -615,9 +616,9 @@ void osd_t::create_osd_state()
// Renew lease
void osd_t::renew_lease(bool reload)
{
st_cli.etcd_call("/lease/keepalive", json11::Json::object {
st_cli->etcd_call("/lease/keepalive", json11::Json::object {
{ "ID", etcd_lease_id }
}, st_cli.etcd_quick_timeout, 0, 0, [this, reload](std::string err, json11::Json data)
}, st_cli->etcd_quick_timeout, 0, 0, [this, reload](std::string err, json11::Json data)
{
if (err == "" && data["result"]["TTL"].uint64_value() == 0)
{
@@ -629,14 +630,14 @@ void osd_t::renew_lease(bool reload)
{
etcd_failed_attempts++;
printf("Error renewing etcd lease: %s\n", err.c_str());
if (etcd_failed_attempts > st_cli.max_etcd_attempts)
if (etcd_failed_attempts > st_cli->max_etcd_attempts)
{
// Die
fprintf(stderr, "Cluster connection failed\n");
force_stop(1);
}
// Retry
tfd->set_timer(st_cli.etcd_quick_timeout, false, [this, reload](int timer_id)
tfd->set_timer(st_cli->etcd_quick_timeout, false, [this, reload](int timer_id)
{
renew_lease(reload);
});
@@ -647,7 +648,7 @@ void osd_t::renew_lease(bool reload)
// Reload PGs
if (reload && run_primary)
{
st_cli.load_pgs();
st_cli->load_pgs();
}
}
});
@@ -657,9 +658,9 @@ void osd_t::force_stop(int exitcode)
{
if (etcd_lease_id != "")
{
st_cli.etcd_call("/kv/lease/revoke", json11::Json::object {
st_cli->etcd_call("/kv/lease/revoke", json11::Json::object {
{ "ID", etcd_lease_id }
}, st_cli.etcd_quick_timeout, st_cli.max_etcd_attempts, 0, [this, exitcode](std::string err, json11::Json data)
}, st_cli->etcd_quick_timeout, st_cli->max_etcd_attempts, 0, [this, exitcode](std::string err, json11::Json data)
{
if (err != "")
{
@@ -682,7 +683,7 @@ json11::Json osd_t::on_load_pgs_checks_hook()
json11::Json::object {
{ "target", "LEASE" },
{ "lease", etcd_lease_id },
{ "key", base64_encode(st_cli.etcd_prefix+"/osd/state/"+std::to_string(osd_num)) },
{ "key", base64_encode(st_cli->etcd_prefix+"/osd/state/"+std::to_string(osd_num)) },
}
};
return checks;
@@ -714,7 +715,7 @@ void osd_t::apply_no_inode_stats()
return;
}
std::vector<uint64_t> no_inode_stats;
for (auto & pool_item: st_cli.pool_config)
for (auto & pool_item: st_cli->pool_config)
{
if (!pool_item.second.used_for_app.empty())
{
@@ -726,7 +727,7 @@ void osd_t::apply_no_inode_stats()
void osd_t::apply_pg_count()
{
for (auto & pool_item: st_cli.pool_config)
for (auto & pool_item: st_cli->pool_config)
{
auto & pool_cfg = pool_item.second;
if (pool_cfg.real_pg_count == 0)
@@ -802,8 +803,8 @@ void osd_t::reshard_continue()
{
again:
auto pool_id = reshard_pools[0];
auto pool_it = st_cli.pool_config.find(pool_id);
if (pool_it == st_cli.pool_config.end() || !pool_it->second.reshard_state)
auto pool_it = st_cli->pool_config.find(pool_id);
if (pool_it == st_cli->pool_config.end() || !pool_it->second.reshard_state)
{
reshard_pools.erase(reshard_pools.begin());
goto again;
@@ -843,7 +844,7 @@ again:
void osd_t::apply_pg_config()
{
bool all_applied = true;
for (auto & pool_item: st_cli.pool_config)
for (auto & pool_item: st_cli->pool_config)
{
auto & pool_cfg = pool_item.second;
if (pool_cfg.reshard_state)
@@ -993,7 +994,7 @@ void osd_t::apply_pg_config()
{
if (pg_osd != this->osd_num && msgr.osd_peers.find(pg_osd) == msgr.osd_peers.end())
{
msgr.connect_peer(pg_osd, st_cli.peer_states[pg_osd]);
msgr.connect_peer(pg_osd, st_cli->peer_states[pg_osd]);
}
}
start_pg_peering(pg);
@@ -1018,7 +1019,7 @@ struct reporting_pg_t
void osd_t::report_pg_states()
{
if (etcd_reporting_pg_state || !this->pg_state_dirty.size() || !st_cli.address_count())
if (etcd_reporting_pg_state || !this->pg_state_dirty.size() || !st_cli->address_count())
{
return;
}
@@ -1035,12 +1036,12 @@ void osd_t::report_pg_states()
}
auto & pg = pg_it->second;
reporting_pgs.push_back((reporting_pg_t){ *it, pg.history_changed });
std::string state_key_base64 = base64_encode(st_cli.etcd_prefix+"/pg/state/"+std::to_string(pg.pool_id)+"/"+std::to_string(pg.pg_num));
std::string state_key_base64 = base64_encode(st_cli->etcd_prefix+"/pg/state/"+std::to_string(pg.pool_id)+"/"+std::to_string(pg.pg_num));
bool pg_state_exists = false;
if (pg.state != PG_STARTING)
{
auto pool_it = st_cli.pool_config.find(pg.pool_id);
if (pool_it != st_cli.pool_config.end())
auto pool_it = st_cli->pool_config.find(pg.pool_id);
if (pool_it != st_cli->pool_config.end())
{
auto pg_it = pool_it->second.pg_config.find(pg.pg_num);
if (pg_it != pool_it->second.pg_config.end() &&
@@ -1054,7 +1055,7 @@ void osd_t::report_pg_states()
{ "target", "MOD" },
{ "key", state_key_base64 },
{ "result", "LESS" },
{ "mod_revision", st_cli.etcd_watch_revision_pg+1 },
{ "mod_revision", st_cli->etcd_watch_revision_pg+1 },
});
continue;
}
@@ -1113,7 +1114,7 @@ void osd_t::report_pg_states()
{
// Prevent race conditions (for the case when the monitor is updating this key at the same time)
pg.history_changed = false;
std::string history_key = base64_encode(st_cli.etcd_prefix+"/pg/history/"+std::to_string(pg.pool_id)+"/"+std::to_string(pg.pg_num));
std::string history_key = base64_encode(st_cli->etcd_prefix+"/pg/history/"+std::to_string(pg.pool_id)+"/"+std::to_string(pg.pg_num));
json11::Json::object history_value = {
{ "epoch", pg.epoch },
{ "all_peers", pg.all_peers },
@@ -1125,7 +1126,7 @@ void osd_t::report_pg_states()
{ "target", "MOD" },
{ "key", history_key },
{ "result", "LESS" },
{ "mod_revision", st_cli.etcd_watch_revision_pg+1 },
{ "mod_revision", st_cli->etcd_watch_revision_pg+1 },
});
success.push_back(json11::Json::object {
{ "request_put", json11::Json::object {
@@ -1143,9 +1144,9 @@ void osd_t::report_pg_states()
}
pg_state_dirty.clear();
etcd_reporting_pg_state = true;
st_cli.etcd_txn(json11::Json::object {
st_cli->etcd_txn(json11::Json::object {
{ "compare", checks }, { "success", success }, { "failure", failure }
}, st_cli.etcd_quick_timeout, 0, 0, [this, reporting_pgs](std::string err, json11::Json data)
}, st_cli->etcd_quick_timeout, 0, 0, [this, reporting_pgs](std::string err, json11::Json data)
{
etcd_reporting_pg_state = false;
if (!data["succeeded"].bool_value())
@@ -1173,13 +1174,13 @@ void osd_t::report_pg_states()
{
if (res["response_range"]["kvs"].array_items().size())
{
auto kv = st_cli.parse_etcd_kv(res["response_range"]["kvs"][0]);
if (kv.key.substr(0, st_cli.etcd_prefix.length()+10) == st_cli.etcd_prefix+"/pg/state/")
auto kv = st_cli->parse_etcd_kv(res["response_range"]["kvs"][0]);
if (kv.key.substr(0, st_cli->etcd_prefix.length()+10) == st_cli->etcd_prefix+"/pg/state/")
{
pool_id_t pool_id = 0;
pg_num_t pg_num = 0;
char null_byte = 0;
int scanned = sscanf(kv.key.c_str() + st_cli.etcd_prefix.length()+10, "%u/%u%c", &pool_id, &pg_num, &null_byte);
int scanned = sscanf(kv.key.c_str() + st_cli->etcd_prefix.length()+10, "%u/%u%c", &pool_id, &pg_num, &null_byte);
if (scanned == 2)
{
auto pg_it = pgs.find({ .pool_id = pool_id, .pg_num = pg_num });
+2 -2
View File
@@ -243,8 +243,8 @@ bool osd_t::pick_next_recovery(osd_recovery_op_t &op)
auto & src = recovery_last_degraded ? pg_it->second.degraded_objects : pg_it->second.misplaced_objects;
if ((pg_it->second.state & mask) == check && src.size() > 0)
{
auto pool_it = st_cli.pool_config.find(pg_it->first.pool_id);
if (pool_it != st_cli.pool_config.end() && pool_it->second.backfillfull)
auto pool_it = st_cli->pool_config.find(pg_it->first.pool_id);
if (pool_it != st_cli->pool_config.end() && pool_it->second.backfillfull)
{
// Skip the pool
recovery_last_pg.pool_id++;
+5 -5
View File
@@ -208,8 +208,8 @@ void osd_t::start_pg_peering(pg_t & pg)
msgr.osd_peers.find(pg_osd) == msgr.osd_peers.end())
{
if (msgr.wanted_peers.find(pg_osd) == msgr.wanted_peers.end())
msgr.connect_peer(pg_osd, st_cli.peer_states[pg_osd]);
if (!st_cli.peer_states[pg_osd].is_null())
msgr.connect_peer(pg_osd, st_cli->peer_states[pg_osd]);
if (!st_cli->peer_states[pg_osd].is_null())
all_connected = false;
}
}
@@ -525,7 +525,7 @@ void osd_t::relock_pg(pg_t & pg)
void osd_t::submit_list_subop(osd_num_t role_osd, pg_peering_state_t *ps)
{
auto & pool_cfg = st_cli.pool_config.at(ps->pool_id);
auto & pool_cfg = st_cli->pool_config.at(ps->pool_id);
if (role_osd == this->osd_num)
{
// Self
@@ -706,7 +706,7 @@ void osd_t::report_pg_state(pg_t & pg)
std::sort(pg.all_peers.begin(), pg.all_peers.end());
pg.cur_peers = pg.target_set;
// Change pg_config at the same time, otherwise our PG reconciling loop may try to apply the old metadata
auto & pg_cfg = st_cli.pool_config[pg.pool_id].pg_config[pg.pg_num];
auto & pg_cfg = st_cli->pool_config[pg.pool_id].pg_config[pg.pg_num];
pg_cfg.target_history = pg.target_history;
pg_cfg.all_peers = pg.all_peers;
}
@@ -748,7 +748,7 @@ void osd_t::report_pg_state(pg_t & pg)
pg.cur_peers.push_back(pg_osd);
}
}
auto & pg_cfg = st_cli.pool_config[pg.pool_id].pg_config[pg.pg_num];
auto & pg_cfg = st_cli->pool_config[pg.pool_id].pg_config[pg.pg_num];
pg_cfg.target_history = pg.target_history;
pg_cfg.all_peers = pg.all_peers;
}
+9 -9
View File
@@ -20,8 +20,8 @@ bool osd_t::prepare_primary_rw(osd_op_t *cur_op)
// K = (pg_size-parity_chunks) in case of EC/XOR, or 1 for replicated pools
pool_id_t pool_id = INODE_POOL(cur_op->req.rw.inode);
// Note: We read pool config here, so we must NOT change it when PGs are active
auto pool_cfg_it = st_cli.pool_config.find(pool_id);
if (pool_cfg_it == st_cli.pool_config.end())
auto pool_cfg_it = st_cli->pool_config.find(pool_id);
if (pool_cfg_it == st_cli->pool_config.end())
{
// Pool config is not loaded yet
finish_op(cur_op, -EPIPE);
@@ -78,7 +78,7 @@ bool osd_t::prepare_primary_rw(osd_op_t *cur_op)
{
// Chained read
// FIXME: Introduce an explicit opcode for chained reads
auto inode_it = st_cli.inode_config.find(cur_op->req.rw.inode);
auto inode_it = st_cli->inode_config.find(cur_op->req.rw.inode);
if (inode_it->second.mod_revision != cur_op->req.rw.meta_revision)
{
// Client view of the metadata differs from OSD's view
@@ -87,14 +87,14 @@ bool osd_t::prepare_primary_rw(osd_op_t *cur_op)
return false;
}
// Find parents from the same pool. Optimized reads only work within pools
while (inode_it != st_cli.inode_config.end() &&
while (inode_it != st_cli->inode_config.end() &&
inode_it->second.parent_id &&
INODE_POOL(inode_it->second.parent_id) == pool_cfg.id)
{
// Check for loops - FIXME check it in etcd_state_client
if (inode_it->second.parent_id == cur_op->req.rw.inode ||
inode_it->second.parent_id == inode_it->second.num ||
chain_size > st_cli.inode_config.size())
chain_size > st_cli->inode_config.size())
{
printf("Inode %ju from pool %u has a parent_id loop, returning EINVAL in response to read\n",
INODE_NO_POOL(cur_op->req.rw.inode), INODE_POOL(cur_op->req.rw.inode));
@@ -102,7 +102,7 @@ bool osd_t::prepare_primary_rw(osd_op_t *cur_op)
return false;
}
chain_size++;
inode_it = st_cli.inode_config.find(inode_it->second.parent_id);
inode_it = st_cli->inode_config.find(inode_it->second.parent_id);
}
if (chain_size)
{
@@ -162,15 +162,15 @@ bool osd_t::prepare_primary_rw(osd_op_t *cur_op)
op_data->read_chain[chain_num] = cur_op->req.rw.inode;
op_data->chain_states[chain_num] = NULL;
chain_num++;
auto inode_it = st_cli.inode_config.find(cur_op->req.rw.inode);
while (inode_it != st_cli.inode_config.end() && inode_it->second.parent_id &&
auto inode_it = st_cli->inode_config.find(cur_op->req.rw.inode);
while (inode_it != st_cli->inode_config.end() && inode_it->second.parent_id &&
INODE_POOL(inode_it->second.parent_id) == pool_cfg.id &&
// Check for loops
inode_it->second.parent_id != cur_op->req.rw.inode)
{
op_data->read_chain[chain_num] = inode_it->second.parent_id;
op_data->chain_states[chain_num] = NULL;
inode_it = st_cli.inode_config.find(inode_it->second.parent_id);
inode_it = st_cli->inode_config.find(inode_it->second.parent_id);
chain_num++;
}
}
+2 -2
View File
@@ -153,14 +153,14 @@ static void add_primary_list(btree::btree_map<object_id, pg_osd_set_state_t*> &
void osd_t::continue_primary_list(osd_op_t *cur_op)
{
auto pool_cfg_it = st_cli.pool_config.find(INODE_POOL(cur_op->req.sec_list.min_inode));
auto pool_cfg_it = st_cli->pool_config.find(INODE_POOL(cur_op->req.sec_list.min_inode));
// Validate the request
if (!cur_op->req.sec_list.list_pg ||
!INODE_POOL(cur_op->req.sec_list.min_inode) ||
INODE_NO_POOL(cur_op->req.sec_list.min_inode) != INODE_NO_POOL(cur_op->req.sec_list.max_inode) ||
INODE_POOL(cur_op->req.sec_list.max_inode) != INODE_POOL(cur_op->req.sec_list.min_inode) ||
cur_op->req.sec_list.stable_limit ||
pool_cfg_it == st_cli.pool_config.end() ||
pool_cfg_it == st_cli->pool_config.end() ||
(cur_op->req.sec_list.pg_stripe_size != 0 && cur_op->req.sec_list.pg_stripe_size != pool_cfg_it->second.pg_stripe_size) ||
(cur_op->req.sec_list.pg_count != 0 && cur_op->req.sec_list.pg_count != pool_cfg_it->second.real_pg_count))
{
+2 -2
View File
@@ -431,9 +431,9 @@ void osd_t::on_change_pg_history_hook(pool_id_t pool_id, pg_num_t pg_num)
}
auto & pg = pg_it->second;
if (pg.epoch > pg.reported_epoch &&
st_cli.pool_config[pool_id].pg_config[pg_num].epoch >= pg.epoch)
st_cli->pool_config[pool_id].pg_config[pg_num].epoch >= pg.epoch)
{
pg.reported_epoch = st_cli.pool_config[pool_id].pg_config[pg_num].epoch;
pg.reported_epoch = st_cli->pool_config[pool_id].pg_config[pg_num].epoch;
std::vector<object_id> resume_oids;
for (auto & op: pg.write_queue)
{
+2 -2
View File
@@ -9,7 +9,7 @@ void osd_t::scrub_list(pool_pg_num_t pg_id, osd_num_t role_osd, object_id min_oi
{
pool_id_t pool_id = pg_id.pool_id;
pg_num_t pg_num = pg_id.pg_num;
auto & pool_cfg = st_cli.pool_config.at(pool_id);
auto & pool_cfg = st_cli->pool_config.at(pool_id);
assert(!scrub_list_op);
if (role_osd == this->osd_num)
{
@@ -327,7 +327,7 @@ void osd_t::plan_scrub(pg_t & pg, bool report_state)
{
timespec tv_now;
clock_gettime(CLOCK_REALTIME, &tv_now);
auto & pool_cfg = st_cli.pool_config.at(pg.pool_id);
auto & pool_cfg = st_cli->pool_config.at(pg.pool_id);
auto interval = pool_cfg.scrub_interval ? pool_cfg.scrub_interval : global_scrub_interval;
if (pg.next_scrub != tv_now.tv_sec + interval)
{
+5 -5
View File
@@ -82,8 +82,8 @@ void osd_t::exec_secondary(osd_op_t *op)
bool osd_t::sec_check_pg_lock(osd_num_t primary_osd, const object_id &oid, uint32_t flags)
{
pool_id_t pool_id = INODE_POOL(oid.inode);
auto pool_cfg_it = st_cli.pool_config.find(pool_id);
if (pool_cfg_it == st_cli.pool_config.end())
auto pool_cfg_it = st_cli->pool_config.find(pool_id);
if (pool_cfg_it == st_cli->pool_config.end())
{
return false;
}
@@ -286,8 +286,8 @@ void osd_t::exec_sec_lock(osd_op_t *cur_op)
return;
}
auto ppg = (pool_pg_num_t){ .pool_id = (pool_id_t)cur_op->req.sec_lock.pool_id, .pg_num = (pg_num_t)cur_op->req.sec_lock.pg_num };
auto pool_cfg_it = st_cli.pool_config.find(ppg.pool_id);
if (pool_cfg_it == st_cli.pool_config.end() ||
auto pool_cfg_it = st_cli->pool_config.find(ppg.pool_id);
if (pool_cfg_it == st_cli->pool_config.end() ||
pool_cfg_it->second.real_pg_count < cur_op->req.sec_lock.pg_num)
{
finish_op(cur_op, -ENOENT);
@@ -354,7 +354,7 @@ void osd_t::exec_show_config(osd_op_t *cur_op)
{ "readonly", readonly },
{ "immediate_commit", (immediate_commit == IMMEDIATE_ALL ? "all" :
(immediate_commit == IMMEDIATE_SMALL ? "small" : "none")) },
{ "lease_timeout", etcd_report_interval+(st_cli.max_etcd_attempts*(2*st_cli.etcd_quick_timeout)+999)/1000 },
{ "lease_timeout", etcd_report_interval+(st_cli->max_etcd_attempts*(2*st_cli->etcd_quick_timeout)+999)/1000 },
{ "features", json11::Json::object{ { "pg_locks", true } } },
};
#ifdef WITH_RDMA