From b2416afb28acb69ff6c6db8437516e2119b98cab Mon Sep 17 00:00:00 2001 From: Vitaliy Filippov Date: Sat, 19 Apr 2025 16:17:38 +0300 Subject: [PATCH] Lock PGs on secondary OSDs to allow local reads and guarantee splitbrain prevention --- docs/config/osd.en.md | 19 +++ docs/config/osd.ru.md | 20 +++ docs/config/src/osd.yml | 20 +++ src/client/messenger.cpp | 15 +- src/client/messenger.h | 2 + src/client/msgr_rdmacm.cpp | 1 + src/client/osd_ops.cpp | 1 + src/client/osd_ops.h | 41 ++++- src/cmd/cli_fix.cpp | 1 + src/osd/osd.cpp | 11 ++ src/osd/osd.h | 20 +++ src/osd/osd_cluster.cpp | 74 ++++++--- src/osd/osd_flush.cpp | 14 +- src/osd/osd_peering.cpp | 268 +++++++++++++++++++++++++++++---- src/osd/osd_peering_pg.cpp | 10 ++ src/osd/osd_peering_pg.h | 10 ++ src/osd/osd_primary.cpp | 2 +- src/osd/osd_primary_subops.cpp | 19 ++- src/osd/osd_primary_sync.cpp | 61 ++++---- src/osd/osd_primary_write.cpp | 82 +++++----- src/osd/osd_secondary.cpp | 160 +++++++++++++++++--- 21 files changed, 676 insertions(+), 175 deletions(-) diff --git a/docs/config/osd.en.md b/docs/config/osd.en.md index 2ae2a302..7baadf1a 100644 --- a/docs/config/osd.en.md +++ b/docs/config/osd.en.md @@ -63,6 +63,8 @@ with an OSD restart or, for some of them, even without restarting by updating co - [discard_on_start](#discard_on_start) - [min_discard_size](#min_discard_size) - [allow_net_split](#allow_net_split) +- [enable_pg_locks](#enable_pg_locks) +- [pg_lock_retry_interval_ms](#pg_lock_retry_interval_ms) ## bind_address @@ -647,3 +649,20 @@ The downside is that it increases the probability of writing data into just pg_m OSDs during failover which can lead to PGs becoming incomplete after additional outages. The old behaviour in versions up to 2.0.0 was equal to enabled allow_net_split. + +## enable_pg_locks + +- Type: boolean + +Vitastor 2.2.0 introduces a new layer of split-brain prevention mechanism in +addition to etcd: PG locks. They prevent split-brain even in abnormal theoretical cases +when etcd is extremely laggy. As a new feature, by default, PG locks are only enabled +for pools where they're required - pools with [localized reads](pool.en.md#local_reads). +Use this parameter to enable or disable this function for all pools. + +## pg_lock_retry_interval_ms + +- Type: milliseconds +- Default: 100 + +Retry interval for failed PG lock attempts. diff --git a/docs/config/osd.ru.md b/docs/config/osd.ru.md index 96c284b7..ead4f369 100644 --- a/docs/config/osd.ru.md +++ b/docs/config/osd.ru.md @@ -64,6 +64,8 @@ - [discard_on_start](#discard_on_start) - [min_discard_size](#min_discard_size) - [allow_net_split](#allow_net_split) +- [enable_pg_locks](#enable_pg_locks) +- [pg_lock_retry_interval_ms](#pg_lock_retry_interval_ms) ## bind_address @@ -679,3 +681,21 @@ pg_minsize OSD во время переключений, что может по неполными (incomplete), если упадут ещё какие-то OSD. Старое поведение в версиях до 2.0.0 было идентично включённому allow_net_split. + +## enable_pg_locks + +- Тип: булево (да/нет) + +В Vitastor 2.2.0 появился новый слой защиты от сплитбрейна в дополнение к etcd - +блокировки PG. Они гарантируют порядок даже в теоретических ненормальных случаях, +когда etcd очень сильно тормозит. Так как функция новая, по умолчанию она включается +только для пулов, в которых она необходима - а именно, в пулах с включёнными +[локальными чтениями](pool.ru.md#local_reads). Ну а с помощью данного параметра +можно включить блокировки PG для всех пулов. + +## pg_lock_retry_interval_ms + +- Тип: миллисекунды +- Значение по умолчанию: 100 + +Интервал повтора неудачных попыток блокировки PG. diff --git a/docs/config/src/osd.yml b/docs/config/src/osd.yml index 1594b4f5..50f1ad58 100644 --- a/docs/config/src/osd.yml +++ b/docs/config/src/osd.yml @@ -781,3 +781,23 @@ неполными (incomplete), если упадут ещё какие-то OSD. Старое поведение в версиях до 2.0.0 было идентично включённому allow_net_split. +- name: enable_pg_locks + type: bool + info: | + Vitastor 2.2.0 introduces a new layer of split-brain prevention mechanism in + addition to etcd: PG locks. They prevent split-brain even in abnormal theoretical cases + when etcd is extremely laggy. As a new feature, by default, PG locks are only enabled + for pools where they're required - pools with [localized reads](pool.en.md#local_reads). + Use this parameter to enable or disable this function for all pools. + info_ru: | + В Vitastor 2.2.0 появился новый слой защиты от сплитбрейна в дополнение к etcd - + блокировки PG. Они гарантируют порядок даже в теоретических ненормальных случаях, + когда etcd очень сильно тормозит. Так как функция новая, по умолчанию она включается + только для пулов, в которых она необходима - а именно, в пулах с включёнными + [локальными чтениями](pool.ru.md#local_reads). Ну а с помощью данного параметра + можно включить блокировки PG для всех пулов. +- name: pg_lock_retry_interval_ms + type: ms + default: 100 + info: Retry interval for failed PG lock attempts. + info_ru: Интервал повтора неудачных попыток блокировки PG. diff --git a/src/client/messenger.cpp b/src/client/messenger.cpp index fed807af..ca53bcf1 100644 --- a/src/client/messenger.cpp +++ b/src/client/messenger.cpp @@ -773,12 +773,15 @@ void osd_messenger_t::accept_connections(int listen_fd) fcntl(peer_fd, F_SETFL, fcntl(peer_fd, F_GETFL, 0) | O_NONBLOCK); int one = 1; setsockopt(peer_fd, SOL_TCP, TCP_NODELAY, &one, sizeof(one)); - clients[peer_fd] = new osd_client_t(); - clients[peer_fd]->peer_addr = addr; - clients[peer_fd]->peer_port = ntohs(((sockaddr_in*)&addr)->sin_port); - clients[peer_fd]->peer_fd = peer_fd; - clients[peer_fd]->peer_state = PEER_CONNECTED; - clients[peer_fd]->in_buf = malloc_or_die(receive_buffer_size); + auto cl = new osd_client_t(); + clients[peer_fd] = cl; + cl->is_incoming = true; + cl->peer_addr = addr; + cl->peer_addr = addr; + cl->peer_port = ntohs(((sockaddr_in*)&addr)->sin_port); + cl->peer_fd = peer_fd; + cl->peer_state = PEER_CONNECTED; + cl->in_buf = malloc_or_die(receive_buffer_size); // Add FD to epoll tfd->set_fd_handler(peer_fd, false, [this](int peer_fd, int epoll_events) { diff --git a/src/client/messenger.h b/src/client/messenger.h index 1754a666..d5d9e277 100644 --- a/src/client/messenger.h +++ b/src/client/messenger.h @@ -60,6 +60,7 @@ struct osd_client_t int ping_time_remaining = 0; int idle_time_remaining = 0; osd_num_t osd_num = 0; + bool is_incoming = false; void *in_buf = NULL; @@ -77,6 +78,7 @@ struct osd_client_t osd_op_buf_list_t recv_list; uint64_t read_op_id = 1; bool check_sequencing = false; + bool enable_pg_locks = false; // Incoming operations std::vector received_ops; diff --git a/src/client/msgr_rdmacm.cpp b/src/client/msgr_rdmacm.cpp index 8be31326..ee4f5d5a 100644 --- a/src/client/msgr_rdmacm.cpp +++ b/src/client/msgr_rdmacm.cpp @@ -510,6 +510,7 @@ void osd_messenger_t::rdmacm_established(rdma_cm_event *ev) rc->qp = conn->cmid->qp; // And an osd_client_t auto cl = new osd_client_t(); + cl->is_incoming = true; cl->peer_addr = conn->parsed_addr; cl->peer_port = conn->rdmacm_port; cl->peer_fd = conn->peer_fd; diff --git a/src/client/osd_ops.cpp b/src/client/osd_ops.cpp index 4cc60f37..cbfaad5b 100644 --- a/src/client/osd_ops.cpp +++ b/src/client/osd_ops.cpp @@ -23,4 +23,5 @@ const char* osd_op_names[] = { "sec_read_bmp", "scrub", "describe", + "sec_lock", }; diff --git a/src/client/osd_ops.h b/src/client/osd_ops.h index 109995be..28dae0a7 100644 --- a/src/client/osd_ops.h +++ b/src/client/osd_ops.h @@ -31,10 +31,13 @@ #define OSD_OP_SEC_READ_BMP 16 #define OSD_OP_SCRUB 17 #define OSD_OP_DESCRIBE 18 -#define OSD_OP_MAX 18 +#define OSD_OP_SEC_LOCK 19 +#define OSD_OP_MAX 19 #define OSD_RW_MAX 64*1024*1024 #define OSD_PROTOCOL_VERSION 1 + #define OSD_OP_RECOVERY_RELATED (uint32_t)1 +#define OSD_OP_IGNORE_PG_LOCK (uint32_t)2 // Memory alignment for direct I/O (usually 512 bytes) #ifndef DIRECT_IO_ALIGNMENT @@ -56,6 +59,9 @@ #define OSD_DEL_SUPPORT_LEFT_ON_DEAD 1 #define OSD_DEL_LEFT_ON_DEAD 2 +#define OSD_SEC_LOCK_PG 1 +#define OSD_SEC_UNLOCK_PG 2 + // common request and reply headers struct __attribute__((__packed__)) osd_op_header_t { @@ -94,7 +100,7 @@ struct __attribute__((__packed__)) osd_op_sec_rw_t uint32_t len; // bitmap/attribute length - bitmap comes after header, but before data uint32_t attr_len; - // the only possible flag is OSD_OP_RECOVERY_RELATED + // OSD_OP_RECOVERY_RELATED, OSD_OP_IGNORE_PG_LOCK uint32_t flags; }; @@ -116,7 +122,7 @@ struct __attribute__((__packed__)) osd_op_sec_del_t object_id oid; // delete version (automatic or specific) uint64_t version; - // the only possible flag is OSD_OP_RECOVERY_RELATED + // OSD_OP_RECOVERY_RELATED, OSD_OP_IGNORE_PG_LOCK uint32_t flags; uint32_t pad0; }; @@ -131,7 +137,7 @@ struct __attribute__((__packed__)) osd_reply_sec_del_t struct __attribute__((__packed__)) osd_op_sec_sync_t { osd_op_header_t header; - // the only possible flag is OSD_OP_RECOVERY_RELATED + // OSD_OP_RECOVERY_RELATED, OSD_OP_IGNORE_PG_LOCK uint32_t flags; uint32_t pad0; }; @@ -147,7 +153,7 @@ struct __attribute__((__packed__)) osd_op_sec_stab_t osd_op_header_t header; // obj_ver_id array length in bytes uint64_t len; - // the only possible flag is OSD_OP_RECOVERY_RELATED + // OSD_OP_RECOVERY_RELATED, OSD_OP_IGNORE_PG_LOCK uint32_t flags; uint32_t pad0; }; @@ -165,6 +171,8 @@ struct __attribute__((__packed__)) osd_op_sec_read_bmp_t osd_op_header_t header; // obj_ver_id array length in bytes uint64_t len; + // OSD_OP_RECOVERY_RELATED, OSD_OP_IGNORE_PG_LOCK + uint32_t flags; }; struct __attribute__((__packed__)) osd_reply_sec_read_bmp_t @@ -173,7 +181,7 @@ struct __attribute__((__packed__)) osd_reply_sec_read_bmp_t osd_reply_header_t header; }; -// show configuration +// show configuration and remember peer information struct __attribute__((__packed__)) osd_op_show_config_t { osd_op_header_t header; @@ -303,6 +311,25 @@ struct __attribute__((__packed__)) osd_reply_describe_item_t osd_num_t osd_num; // OSD number }; +// lock/unlock PG for use by a primary OSD +struct __attribute__((__packed__)) osd_op_sec_lock_t +{ + osd_op_header_t header; + // OSD_SEC_LOCK_PG or OSD_SEC_UNLOCK_PG + uint64_t flags; + // Pool ID and PG number + uint64_t pool_id; + uint64_t pg_num; + // PG state as calculated by the primary OSD + uint64_t pg_state; +}; + +struct __attribute__((__packed__)) osd_reply_sec_lock_t +{ + osd_reply_header_t header; + uint64_t cur_primary; +}; + // FIXME it would be interesting to try to unify blockstore_op and osd_op formats union osd_any_op_t { @@ -313,6 +340,7 @@ union osd_any_op_t osd_op_sec_stab_t sec_stab; osd_op_sec_read_bmp_t sec_read_bmp; osd_op_sec_list_t sec_list; + osd_op_sec_lock_t sec_lock; osd_op_show_config_t show_conf; osd_op_rw_t rw; osd_op_sync_t sync; @@ -329,6 +357,7 @@ union osd_any_reply_t osd_reply_sec_stab_t sec_stab; osd_reply_sec_read_bmp_t sec_read_bmp; osd_reply_sec_list_t sec_list; + osd_reply_sec_lock_t sec_lock; osd_reply_show_config_t show_conf; osd_reply_rw_t rw; osd_reply_del_t del; diff --git a/src/cmd/cli_fix.cpp b/src/cmd/cli_fix.cpp index 6f57dc4c..d930fd9d 100644 --- a/src/cmd/cli_fix.cpp +++ b/src/cmd/cli_fix.cpp @@ -200,6 +200,7 @@ struct cli_fix_t .stripe = op->req.describe.min_offset | items[i].role, }, .version = 0, + .flags = OSD_OP_IGNORE_PG_LOCK, }, }; rm_op->callback = [this, primary_osd, rm_osd_num, rm_count, &obj](osd_op_t *rm_op) diff --git a/src/osd/osd.cpp b/src/osd/osd.cpp index 667e8cad..d44094e0 100644 --- a/src/osd/osd.cpp +++ b/src/osd/osd.cpp @@ -271,6 +271,12 @@ void osd_t::parse_config(bool init) inode_vanish_time = config["inode_vanish_time"].uint64_value(); if (!inode_vanish_time) inode_vanish_time = 60; + enable_pg_locks = config["enable_pg_locks"].is_null() || json_is_true(config["enable_pg_locks"]); + bool old_pg_locks_localize_only = pg_locks_localize_only; + pg_locks_localize_only = config["enable_pg_locks"].is_null(); + pg_lock_retry_interval_ms = config["pg_lock_retry_interval"].uint64_value(); + if (pg_lock_retry_interval_ms <= 1) + pg_lock_retry_interval_ms = 100; auto old_auto_scrub = auto_scrub; auto_scrub = json_is_true(config["auto_scrub"]); global_scrub_interval = parse_time(config["scrub_interval"].string_value()); @@ -336,6 +342,10 @@ void osd_t::parse_config(bool init) { apply_recovery_tune_interval(); } + if (old_pg_locks_localize_only != pg_locks_localize_only) + { + apply_pg_locks_localize_only(); + } } void osd_t::bind_socket() @@ -447,6 +457,7 @@ void osd_t::exec_op(osd_op_t *cur_op) } if (readonly && cur_op->req.hdr.opcode != OSD_OP_SEC_READ && + cur_op->req.hdr.opcode != OSD_OP_SEC_LOCK && cur_op->req.hdr.opcode != OSD_OP_SEC_LIST && cur_op->req.hdr.opcode != OSD_OP_READ && cur_op->req.hdr.opcode != OSD_OP_SEC_READ_BMP && diff --git a/src/osd/osd.h b/src/osd/osd.h index 576fa85a..79afc866 100644 --- a/src/osd/osd.h +++ b/src/osd/osd.h @@ -92,6 +92,12 @@ struct recovery_stat_t uint64_t count, usec, bytes; }; +struct osd_pg_lock_t +{ + osd_num_t primary_osd = 0; + uint64_t state = 0; +}; + class osd_t { // config @@ -140,6 +146,9 @@ class osd_t uint32_t scrub_list_limit = 1000; bool scrub_find_best = true; uint64_t scrub_ec_max_bruteforce = 100; + bool enable_pg_locks = false; + bool pg_locks_localize_only = false; + uint64_t pg_lock_retry_interval_ms = 100; // cluster state @@ -159,6 +168,7 @@ class osd_t // peers and PGs + std::map pg_locks; std::map pg_counts; std::map pgs; std::set dirty_pgs; @@ -239,6 +249,8 @@ class osd_t void on_change_etcd_state_hook(std::map & changes); void on_load_config_hook(json11::Json::object & changes); void on_reload_config_hook(json11::Json::object & changes); + void on_change_pool_config_hook(); + void apply_pg_locks_localize_only(); json11::Json on_load_pgs_checks_hook(); void on_load_pgs_hook(bool success); void bind_socket(); @@ -268,11 +280,16 @@ class osd_t void repeer_pgs(osd_num_t osd_num); void start_pg_peering(pg_t & pg); void drop_dirty_pg_connections(pool_pg_num_t pg); + void record_pg_lock(pg_t & pg, osd_num_t peer_osd, uint64_t pg_state); + void relock_pg(pg_t & pg); void submit_list_subop(osd_num_t role_osd, pg_peering_state_t *ps); void discard_list_subop(osd_op_t *list_op); bool stop_pg(pg_t & pg); void reset_pg(pg_t & pg); void finish_stop_pg(pg_t & pg); + void rm_inflight(pg_t & pg); + void continue_pg(pg_t & pg); + bool continue_pg_peering(pg_t & pg); // flushing, recovery and backfill void submit_pg_flush_ops(pg_t & pg); @@ -299,10 +316,13 @@ class osd_t void finish_op(osd_op_t *cur_op, int retval); // secondary ops + bool sec_check_pg_lock(osd_num_t primary_osd, const object_id &oid); void exec_sync_stab_all(osd_op_t *cur_op); void exec_show_config(osd_op_t *cur_op); void exec_secondary(osd_op_t *cur_op); void exec_secondary_real(osd_op_t *cur_op); + void exec_sec_read_bmp(osd_op_t *cur_op); + void exec_sec_lock(osd_op_t *cur_op); void secondary_op_callback(osd_op_t *cur_op); // primary ops diff --git a/src/osd/osd_cluster.cpp b/src/osd/osd_cluster.cpp index 39b4f593..4b3ea2df 100644 --- a/src/osd/osd_cluster.cpp +++ b/src/osd/osd_cluster.cpp @@ -65,6 +65,7 @@ void osd_t::init_cluster() 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 & changes) { on_change_etcd_state_hook(changes); }; @@ -153,6 +154,7 @@ bool osd_t::check_peer_config(osd_client_t *cl, json11::Json conf) return false; } } + cl->enable_pg_locks = conf["features"]["pg_locks"].bool_value(); return true; } @@ -414,6 +416,28 @@ void osd_t::on_change_osd_state_hook(osd_num_t peer_osd) } } +void osd_t::on_change_pool_config_hook() +{ + apply_pg_locks_localize_only(); +} + +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()) + { + continue; + } + auto & pool_cfg = pool_it->second; + auto & pg = pp.second; + pg.disable_pg_locks = pg_locks_localize_only && + pool_cfg.scheme == POOL_SCHEME_REPLICATED && + pool_cfg.local_reads == POOL_LOCAL_READ_PRIMARY; + } +} + void osd_t::on_change_backfillfull_hook(pool_id_t pool_id) { if (!(peering_state & (OSD_RECOVERING | OSD_FLUSHING_PGS))) @@ -691,20 +715,27 @@ void osd_t::apply_pg_count() // The external tool must wait for all PGs to come down before changing PG count // If it doesn't wait, a restarted OSD may apply the new count immediately which will lead to bugs // So an OSD just dies if it detects PG count change while there are active PGs - int still_active = 0; + int still_active_primary = 0; for (auto & kv: pgs) { if (kv.first.pool_id == pool_item.first && (kv.second.state & PG_ACTIVE)) { - still_active++; + still_active_primary++; } } - if (still_active > 0) + int still_active_secondary = 0; + for (auto lock_it = pg_locks.lower_bound((pool_pg_num_t){ .pool_id = pool_item.first, .pg_num = 0 }); + lock_it != pg_locks.end() && lock_it->first.pool_id == pool_item.first; lock_it++) + { + still_active_secondary++; + } + if (still_active_primary > 0 || still_active_secondary > 0) { printf( "[OSD %ju] PG count change detected for pool %u (new is %ju, old is %u)," - " but %u PG(s) are still active. This is not allowed. Exiting\n", - this->osd_num, pool_item.first, pool_item.second.real_pg_count, pg_counts[pool_item.first], still_active + " but %u PG(s) are still active as primary and %u as secondary. This is not allowed. Exiting\n", + this->osd_num, pool_item.first, pool_item.second.real_pg_count, pg_counts[pool_item.first], + still_active_primary, still_active_secondary ); force_stop(1); return; @@ -831,22 +862,23 @@ void osd_t::apply_pg_config() } } auto & pg = this->pgs[{ .pool_id = pool_id, .pg_num = pg_num }]; - pg = (pg_t){ - .state = pg_cfg.cur_primary == this->osd_num ? PG_PEERING : PG_STARTING, - .scheme = pool_item.second.scheme, - .pg_cursize = 0, - .pg_size = pool_item.second.pg_size, - .pg_minsize = pool_item.second.pg_minsize, - .pg_data_size = pool_item.second.scheme == POOL_SCHEME_REPLICATED - ? 1 : pool_item.second.pg_size - pool_item.second.parity_chunks, - .pool_id = pool_id, - .pg_num = pg_num, - .reported_epoch = pg_cfg.epoch, - .target_history = pg_cfg.target_history, - .all_peers = vec_all_peers, - .next_scrub = pg_cfg.next_scrub, - .target_set = pg_cfg.target_set, - }; + pg.state = pg_cfg.cur_primary == this->osd_num ? PG_PEERING : PG_STARTING; + pg.scheme = pool_item.second.scheme; + pg.pg_cursize = 0; + pg.pg_size = pool_item.second.pg_size; + pg.pg_minsize = pool_item.second.pg_minsize; + pg.pg_data_size = pool_item.second.scheme == POOL_SCHEME_REPLICATED + ? 1 : pool_item.second.pg_size - pool_item.second.parity_chunks; + pg.pool_id = pool_id; + pg.pg_num = pg_num; + pg.reported_epoch = pg_cfg.epoch; + pg.target_history = pg_cfg.target_history; + pg.all_peers = vec_all_peers; + pg.next_scrub = pg_cfg.next_scrub; + pg.target_set = pg_cfg.target_set; + pg.disable_pg_locks = pg_locks_localize_only && + pool_item.second.scheme == POOL_SCHEME_REPLICATED && + pool_item.second.local_reads == POOL_LOCAL_READ_PRIMARY; if (pg.scheme == POOL_SCHEME_EC) { use_ec(pg.pg_size, pg.pg_data_size, true); diff --git a/src/osd/osd_flush.cpp b/src/osd/osd_flush.cpp index a663188e..8a63fcee 100644 --- a/src/osd/osd_flush.cpp +++ b/src/osd/osd_flush.cpp @@ -150,14 +150,7 @@ void osd_t::handle_flush_op(bool rollback, pool_id_t pool_id, pg_num_t pg_num, p { continue_primary_write(op); } - if ((pg.state & PG_STOPPING) && pg.inflight == 0 && !pg.flush_batch) - { - finish_stop_pg(pg); - } - else if ((pg.state & PG_REPEERING) && pg.inflight == 0 && !pg.flush_batch) - { - start_pg_peering(pg); - } + continue_pg(pg); } } @@ -254,7 +247,8 @@ bool osd_t::pick_next_recovery(osd_recovery_op_t &op) restart: for (auto pg_it = pgs.lower_bound(recovery_last_pg); pg_it != pgs.end(); pg_it++) { - if ((pg_it->second.state & mask) == check) + 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) @@ -263,8 +257,6 @@ bool osd_t::pick_next_recovery(osd_recovery_op_t &op) recovery_last_pg.pool_id++; goto restart; } - auto & src = recovery_last_degraded ? pg_it->second.degraded_objects : pg_it->second.misplaced_objects; - assert(src.size() > 0); // Restart scanning from the next object for (auto obj_it = src.upper_bound(recovery_last_oid); obj_it != src.end(); obj_it++) { diff --git a/src/osd/osd_peering.cpp b/src/osd/osd_peering.cpp index 5477a5cb..724b5811 100644 --- a/src/osd/osd_peering.cpp +++ b/src/osd/osd_peering.cpp @@ -21,28 +21,8 @@ void osd_t::handle_peers() { if (p.second.state == PG_PEERING) { - if (!p.second.peering_state->list_ops.size()) + if (continue_pg_peering(p.second)) { - p.second.calc_object_states(log_level); - report_pg_state(p.second); - schedule_scrub(p.second); - incomplete_objects += p.second.incomplete_objects.size(); - misplaced_objects += p.second.misplaced_objects.size(); - // FIXME: degraded objects may currently include misplaced, too! Report them separately? - degraded_objects += p.second.degraded_objects.size(); - if (p.second.state & PG_HAS_UNCLEAN) - peering_state = peering_state | OSD_FLUSHING_PGS; - else if (p.second.state & (PG_HAS_DEGRADED | PG_HAS_MISPLACED)) - { - peering_state = peering_state | OSD_RECOVERING; - if (p.second.state & PG_HAS_DEGRADED) - { - // Restart recovery from degraded objects - recovery_last_degraded = true; - recovery_last_pg = {}; - recovery_last_oid = {}; - } - } ringloop->wakeup(); return; } @@ -95,6 +75,16 @@ void osd_t::handle_peers() void osd_t::repeer_pgs(osd_num_t peer_osd) { + if (msgr.osd_peer_fds.find(peer_osd) == msgr.osd_peer_fds.end()) + { + for (auto lock_it = pg_locks.begin(); lock_it != pg_locks.end(); ) + { + if (lock_it->second.primary_osd == peer_osd) + pg_locks.erase(lock_it++); + else + lock_it++; + } + } // Re-peer affected PGs for (auto & p: pgs) { @@ -114,7 +104,7 @@ void osd_t::repeer_pgs(osd_num_t peer_osd) { // Repeer this pg printf("[PG %u/%u] Repeer because of OSD %ju\n", pg.pool_id, pg.pg_num, peer_osd); - if (!(pg.state & (PG_ACTIVE | PG_REPEERING)) || pg.inflight == 0 && !pg.flush_batch) + if (!(pg.state & (PG_ACTIVE | PG_REPEERING)) || pg.can_repeer()) { start_pg_peering(pg); } @@ -195,7 +185,6 @@ void osd_t::start_pg_peering(pg_t & pg) pg.state = PG_PEERING; this->peering_state |= OSD_PEERING_PGS; reset_pg(pg); - report_pg_state(pg); drop_dirty_pg_connections({ .pool_id = pg.pool_id, .pg_num = pg.pg_num }); // Try to connect with current peers if they're up, but we don't have connections to them // Otherwise we may erroneously decide that the pg is incomplete :-) @@ -322,16 +311,203 @@ void osd_t::start_pg_peering(pg_t & pg) pg.peering_state->pool_id = pg.pool_id; pg.peering_state->pg_num = pg.pg_num; } - for (osd_num_t peer_osd: cur_peers) + pg.peering_state->locked = false; + pg.peering_state->lists_done = false; + report_pg_state(pg); +} + +bool osd_t::continue_pg_peering(pg_t & pg) +{ + if (pg.peering_state->locked) { - if (pg.peering_state->list_ops.find(peer_osd) != pg.peering_state->list_ops.end() || - pg.peering_state->list_results.find(peer_osd) != pg.peering_state->list_results.end()) + pg.peering_state->lists_done = true; + for (osd_num_t peer_osd: pg.cur_peers) { + if (pg.peering_state->list_results.find(peer_osd) == pg.peering_state->list_results.end()) + { + pg.peering_state->lists_done = false; + } + if (pg.peering_state->list_ops.find(peer_osd) != pg.peering_state->list_ops.end() || + pg.peering_state->list_results.find(peer_osd) != pg.peering_state->list_results.end()) + { + continue; + } + submit_list_subop(peer_osd, pg.peering_state); + } + } + if (pg.peering_state->lists_done) + { + pg.calc_object_states(log_level); + report_pg_state(pg); + schedule_scrub(pg); + incomplete_objects += pg.incomplete_objects.size(); + misplaced_objects += pg.misplaced_objects.size(); + // FIXME: degraded objects may currently include misplaced, too! Report them separately? + degraded_objects += pg.degraded_objects.size(); + if (pg.state & PG_HAS_UNCLEAN) + this->peering_state = peering_state | OSD_FLUSHING_PGS; + else if (pg.state & (PG_HAS_DEGRADED | PG_HAS_MISPLACED)) + { + this->peering_state = peering_state | OSD_RECOVERING; + if (pg.state & PG_HAS_DEGRADED) + { + // Restart recovery from degraded objects + this->recovery_last_degraded = true; + this->recovery_last_pg = {}; + this->recovery_last_oid = {}; + } + } + return true; + } + return false; +} + +void osd_t::record_pg_lock(pg_t & pg, osd_num_t peer_osd, uint64_t pg_state) +{ + if (!pg_state) + pg.lock_peers.erase(peer_osd); + else + pg.lock_peers[peer_osd] = pg_state; +} + +void osd_t::relock_pg(pg_t & pg) +{ + if (!enable_pg_locks || pg.disable_pg_locks && !pg.lock_peers.size()) + { + if (pg.state & PG_PEERING) + pg.peering_state->locked = true; + continue_pg(pg); + return; + } + if (pg.inflight_locks > 0 || pg.lock_waiting) + { + return; + } + // Check that lock_peers are equal to cur_peers and correct the difference, if any + uint64_t wanted_state = pg.state; + std::vector diff_osds; + if (!(pg.state & (PG_STOPPING | PG_OFFLINE | PG_INCOMPLETE)) && !pg.disable_pg_locks) + { + for (osd_num_t peer_osd: pg.cur_peers) + { + if (peer_osd != this->osd_num) + { + auto lock_it = pg.lock_peers.find(peer_osd); + if (lock_it == pg.lock_peers.end()) + diff_osds.push_back(peer_osd); + else + { + if (lock_it->second != wanted_state) + diff_osds.push_back(peer_osd); + lock_it->second |= ((uint64_t)1 << 63); + } + } + } + } + int relock_osd_count = diff_osds.size(); + for (auto & lp: pg.lock_peers) + { + if (!(lp.second & ((uint64_t)1 << 63))) + diff_osds.push_back(lp.first); + lp.second &= ~((uint64_t)1 << 63); + } + if (!diff_osds.size()) + { + if (pg.state & PG_PEERING) + pg.peering_state->locked = true; + continue_pg(pg); + return; + } + pg.inflight_locks++; + for (int i = 0; i < diff_osds.size(); i++) + { + bool unlock_peer = (i >= relock_osd_count); + uint64_t new_state = unlock_peer ? 0 : pg.state; + auto peer_osd = diff_osds[i]; + auto peer_fd_it = msgr.osd_peer_fds.find(peer_osd); + if (peer_fd_it == msgr.osd_peer_fds.end()) + { + if (unlock_peer) + { + // Peer is dead - unlocked automatically + record_pg_lock(pg, peer_osd, new_state); + diff_osds.erase(diff_osds.begin()+(i--)); + } continue; } - submit_list_subop(peer_osd, pg.peering_state); + int peer_fd = peer_fd_it->second; + auto cl = msgr.clients.at(peer_fd); + if (!cl->enable_pg_locks) + { + // Peer does not support locking - just instantly remember the lock as successful + record_pg_lock(pg, peer_osd, new_state); + diff_osds.erase(diff_osds.begin()+(i--)); + continue; + } + pg.inflight_locks++; + osd_op_t *op = new osd_op_t(); + op->op_type = OSD_OP_OUT; + op->peer_fd = peer_fd; + op->req = (osd_any_op_t){ + .sec_lock = { + .header = { + .magic = SECONDARY_OSD_OP_MAGIC, + .opcode = OSD_OP_SEC_LOCK, + }, + .flags = (uint64_t)(unlock_peer ? OSD_SEC_UNLOCK_PG : OSD_SEC_LOCK_PG), + .pool_id = pg.pool_id, + .pg_num = pg.pg_num, + .pg_state = new_state, + }, + }; + op->callback = [this, peer_osd](osd_op_t *op) + { + pool_pg_num_t pg_id = { .pool_id = (pool_id_t)op->req.sec_lock.pool_id, .pg_num = (pg_num_t)op->req.sec_lock.pg_num }; + auto pg_it = pgs.find(pg_id); + if (pg_it == pgs.end()) + { + return; + } + auto & pg = pg_it->second; + if (op->reply.hdr.retval == 0) + { + record_pg_lock(pg_it->second, peer_osd, op->req.sec_lock.pg_state); + } + else if (op->reply.hdr.retval != -EPIPE) + { + printf( + (op->reply.hdr.retval == -ENOENT + ? "Failed to %1$s PG %2$u/%3$u on OSD %4$ju - peer didn't load PG info yet\n" + : (op->reply.sec_lock.cur_primary + ? "Failed to %1$s PG %2$u/%3$u on OSD %4$ju - taken by OSD %6$ju (retval=%5$jd)\n" + : "Failed to %1$s PG %2$u/%3$u on OSD %4$ju - retval=%5$jd\n")), + op->req.sec_lock.flags == OSD_SEC_UNLOCK_PG ? "unlock" : "lock", + pg_id.pool_id, pg_id.pg_num, peer_osd, op->reply.hdr.retval, op->reply.sec_lock.cur_primary + ); + // Retry relocking/unlocking PG after a short time + pg.lock_waiting = true; + tfd->set_timer(pg_lock_retry_interval_ms, false, [this, pg_id](int) + { + auto pg_it = pgs.find(pg_id); + if (pg_it != pgs.end()) + { + pg_it->second.lock_waiting = false; + relock_pg(pg_it->second); + } + }); + } + pg.inflight_locks--; + relock_pg(pg); + delete op; + }; + msgr.outbox_push(op); } - ringloop->wakeup(); + if (pg.state & PG_PEERING) + { + pg.peering_state->locked = !diff_osds.size(); + } + pg.inflight_locks--; + continue_pg(pg); } void osd_t::submit_list_subop(osd_num_t role_osd, pg_peering_state_t *ps) @@ -379,10 +555,16 @@ void osd_t::submit_list_subop(osd_num_t role_osd, pg_peering_state_t *ps) } else { + auto role_fd_it = msgr.osd_peer_fds.find(role_osd); + if (role_fd_it == msgr.osd_peer_fds.end()) + { + printf("Failed to get object list from OSD %ju because it is disconnected\n", role_osd); + return; + } // Peer osd_op_t *op = new osd_op_t(); op->op_type = OSD_OP_OUT; - op->peer_fd = msgr.osd_peer_fds.at(role_osd); + op->peer_fd = role_fd_it->second; op->req = (osd_any_op_t){ .sec_list = { .header = { @@ -474,8 +656,8 @@ bool osd_t::stop_pg(pg_t & pg) return false; } drop_dirty_pg_connections({ .pool_id = pg.pool_id, .pg_num = pg.pg_num }); - pg.state = pg.state & ~PG_ACTIVE & ~PG_REPEERING | PG_STOPPING; - if (pg.inflight == 0 && !pg.flush_batch) + pg.state = pg.state & ~PG_STARTING & ~PG_PEERING & ~PG_INCOMPLETE & ~PG_ACTIVE & ~PG_REPEERING & ~PG_OFFLINE | PG_STOPPING; + if (pg.can_stop()) { finish_stop_pg(pg); } @@ -556,9 +738,33 @@ void osd_t::report_pg_state(pg_t & pg) pg_cfg.target_history = pg.target_history; pg_cfg.all_peers = pg.all_peers; } + relock_pg(pg); if (pg.state == PG_OFFLINE && !this->pg_config_applied) { apply_pg_config(); } report_pg_states(); } + +void osd_t::rm_inflight(pg_t & pg) +{ + pg.inflight--; + assert(pg.inflight >= 0); + continue_pg(pg); +} + +void osd_t::continue_pg(pg_t & pg) +{ + if ((pg.state & PG_STOPPING) && pg.can_stop()) + { + finish_stop_pg(pg); + } + else if ((pg.state & PG_REPEERING) && pg.can_repeer()) + { + start_pg_peering(pg); + } + else if ((pg.state & PG_PEERING) && pg.peering_state->locked) + { + continue_pg_peering(pg); + } +} diff --git a/src/osd/osd_peering_pg.cpp b/src/osd/osd_peering_pg.cpp index 7dbe6d8b..67895073 100644 --- a/src/osd/osd_peering_pg.cpp +++ b/src/osd/osd_peering_pg.cpp @@ -489,3 +489,13 @@ void pg_t::print_state() total_count ); } + +bool pg_t::can_stop() +{ + return inflight == 0 && inflight_locks == 0 && !lock_peers.size() && !flush_batch; +} + +bool pg_t::can_repeer() +{ + return inflight == 0 && !flush_batch; +} diff --git a/src/osd/osd_peering_pg.h b/src/osd/osd_peering_pg.h index 593f43cb..f90609be 100644 --- a/src/osd/osd_peering_pg.h +++ b/src/osd/osd_peering_pg.h @@ -49,6 +49,8 @@ struct pg_peering_state_t std::map list_results; pool_id_t pool_id = 0; pg_num_t pg_num = 0; + bool locked = false; + bool lists_done = false; }; struct obj_piece_id_t @@ -87,6 +89,7 @@ struct pg_t pool_id_t pool_id = 0; pg_num_t pg_num = 0; uint64_t clean_count = 0, total_count = 0; + bool disable_pg_locks = false; // epoch number - should increase with each non-clean activation of the PG uint64_t epoch = 0, reported_epoch = 0; // target history and all potential peers @@ -104,6 +107,10 @@ struct pg_t // cur_set is the current set of connected peer OSDs for this PG // cur_set = (role => osd_num or UINT64_MAX if missing). role numbers begin with zero std::vector cur_set; + // locked peer list => pg state reported to the peer + std::map lock_peers; + int inflight_locks = 0; + bool lock_waiting = false; // same thing in state_dict-like format pg_osd_set_t cur_loc_set; // moved object map. by default, each object is considered to reside on cur_set. @@ -125,6 +132,9 @@ struct pg_t pg_osd_set_state_t* add_object_to_state(const object_id oid, const uint64_t state, const pg_osd_set_t & osd_set); void calc_object_states(int log_level); void print_state(); + bool can_stop(); + bool can_repeer(); + void rm_inflight(); }; inline bool operator < (const pg_obj_loc_t &a, const pg_obj_loc_t &b) diff --git a/src/osd/osd_primary.cpp b/src/osd/osd_primary.cpp index 0c36afa5..244650e9 100644 --- a/src/osd/osd_primary.cpp +++ b/src/osd/osd_primary.cpp @@ -632,7 +632,7 @@ void osd_t::remove_object_from_state(object_id & oid, pg_osd_set_state_t **objec { this->misplaced_objects--; pg.misplaced_objects.erase(oid); - if (!pg.misplaced_objects.size()) + if (!pg.misplaced_objects.size() && !pg.copies_to_delete_after_sync.size()) { pg.state = pg.state & ~PG_HAS_MISPLACED; changed = true; diff --git a/src/osd/osd_primary_subops.cpp b/src/osd/osd_primary_subops.cpp index 3549d078..85df4bc6 100644 --- a/src/osd/osd_primary_subops.cpp +++ b/src/osd/osd_primary_subops.cpp @@ -77,16 +77,7 @@ void osd_t::finish_op(osd_op_t *cur_op, int retval) if (cur_op->op_data->pg_num > 0) { auto & pg = *cur_op->op_data->pg; - pg.inflight--; - assert(pg.inflight >= 0); - if ((pg.state & PG_STOPPING) && pg.inflight == 0 && !pg.flush_batch) - { - finish_stop_pg(pg); - } - else if ((pg.state & PG_REPEERING) && pg.inflight == 0 && !pg.flush_batch) - { - start_pg_peering(pg); - } + rm_inflight(pg); } assert(!cur_op->op_data->subops); free(cur_op->op_data); @@ -434,6 +425,14 @@ void osd_t::handle_primary_subop(osd_op_t *subop, osd_op_t *cur_op) retval, expected, peer_osd ); } + else if (opcode == OSD_OP_SEC_DELETE) + { + printf( + "delete subop to %jx:%jx v%ju failed on osd %jd: retval = %d (expected %d)\n", + subop->req.sec_del.oid.inode, subop->req.sec_del.oid.stripe, subop->req.sec_del.version, + peer_osd, retval, expected + ); + } else { printf( diff --git a/src/osd/osd_primary_sync.cpp b/src/osd/osd_primary_sync.cpp index 0808773b..aecffa90 100644 --- a/src/osd/osd_primary_sync.cpp +++ b/src/osd/osd_primary_sync.cpp @@ -80,15 +80,17 @@ resume_2: this->unstable_writes.clear(); } { + op_data->dirty_pg_count = dirty_pgs.size(); + op_data->dirty_osd_count = dirty_osds.size(); void *dirty_buf = malloc_or_die( sizeof(pool_pg_num_t)*dirty_pgs.size() + + sizeof(uint64_t)*dirty_pgs.size() + sizeof(osd_num_t)*dirty_osds.size() + sizeof(obj_ver_osd_t)*this->copies_to_delete_after_sync_count ); op_data->dirty_pgs = (pool_pg_num_t*)dirty_buf; - op_data->dirty_osds = (osd_num_t*)((uint8_t*)dirty_buf + sizeof(pool_pg_num_t)*dirty_pgs.size()); - op_data->dirty_pg_count = dirty_pgs.size(); - op_data->dirty_osd_count = dirty_osds.size(); + uint64_t *pg_del_counts = (uint64_t*)((uint8_t*)op_data->dirty_pgs + (sizeof(pool_pg_num_t))*op_data->dirty_pg_count); + op_data->dirty_osds = (osd_num_t*)((uint8_t*)pg_del_counts + 8*op_data->dirty_pg_count); if (this->copies_to_delete_after_sync_count) { op_data->copies_to_delete_count = 0; @@ -103,16 +105,16 @@ resume_2: sizeof(obj_ver_osd_t)*pg.copies_to_delete_after_sync.size() ); op_data->copies_to_delete_count += pg.copies_to_delete_after_sync.size(); - this->copies_to_delete_after_sync_count -= pg.copies_to_delete_after_sync.size(); - pg.copies_to_delete_after_sync.clear(); } - assert(this->copies_to_delete_after_sync_count == 0); } int dpg = 0; for (auto dirty_pg_num: dirty_pgs) { - pgs.at(dirty_pg_num).inflight++; - op_data->dirty_pgs[dpg++] = dirty_pg_num; + auto & pg = pgs.at(dirty_pg_num); + pg.inflight++; + op_data->dirty_pgs[dpg] = dirty_pg_num; + pg_del_counts[dpg] = pg.copies_to_delete_after_sync.size(); + dpg++; } dirty_pgs.clear(); dpg = 0; @@ -183,23 +185,6 @@ resume_6: } } } - if (op_data->copies_to_delete) - { - // Return 'copies to delete' back into respective PGs - for (int i = 0; i < op_data->copies_to_delete_count; i++) - { - auto & w = op_data->copies_to_delete[i]; - auto & pg = pgs.at((pool_pg_num_t){ - .pool_id = INODE_POOL(w.oid.inode), - .pg_num = map_to_pg(w.oid, st_cli.pool_config.at(INODE_POOL(w.oid.inode)).pg_stripe_size), - }); - if (pg.state & PG_ACTIVE) - { - pg.copies_to_delete_after_sync.push_back(w); - copies_to_delete_after_sync_count++; - } - } - } } else if (op_data->copies_to_delete) { @@ -213,6 +198,22 @@ resume_8: { goto resume_6; } + { + uint64_t *pg_del_counts = (uint64_t*)((uint8_t*)op_data->dirty_pgs + (sizeof(pool_pg_num_t))*op_data->dirty_pg_count); + for (int i = 0; i < op_data->dirty_pg_count; i++) + { + auto & pg = pgs.at(op_data->dirty_pgs[i]); + auto n = pg_del_counts[i]; + assert(copies_to_delete_after_sync_count >= n); + copies_to_delete_after_sync_count -= n; + pg.copies_to_delete_after_sync.erase(pg.copies_to_delete_after_sync.begin(), pg.copies_to_delete_after_sync.begin()+n); + if (!pg.misplaced_objects.size() && !pg.copies_to_delete_after_sync.size() && (pg.state & PG_HAS_MISPLACED)) + { + pg.state = pg.state & ~PG_HAS_MISPLACED; + report_pg_state(pg); + } + } + } if (immediate_commit == IMMEDIATE_NONE) { // Mark OSDs as dirty because deletions have to be synced too! @@ -226,15 +227,7 @@ resume_8: for (int i = 0; i < op_data->dirty_pg_count; i++) { auto & pg = pgs.at(op_data->dirty_pgs[i]); - pg.inflight--; - if ((pg.state & PG_STOPPING) && pg.inflight == 0 && !pg.flush_batch) - { - finish_stop_pg(pg); - } - else if ((pg.state & PG_REPEERING) && pg.inflight == 0 && !pg.flush_batch) - { - start_pg_peering(pg); - } + rm_inflight(pg); } // FIXME: Free those in the destructor (not here)? free(op_data->dirty_pgs); diff --git a/src/osd/osd_primary_write.cpp b/src/osd/osd_primary_write.cpp index e00a8a9d..d6c0ba8b 100644 --- a/src/osd/osd_primary_write.cpp +++ b/src/osd/osd_primary_write.cpp @@ -301,6 +301,38 @@ resume_12: } if (op_data->object_state) { + // Any kind of a non-clean object can have extra chunks, because we don't record objects + // as degraded & misplaced or incomplete & misplaced at the same time. So try to remove extra chunks + if (immediate_commit != IMMEDIATE_ALL) + { + // We can't remove extra chunks yet if fsyncs are explicit, because + // new copies may not be committed to stable storage yet + // We can only remove extra chunks after a successful SYNC for this PG + for (auto & chunk: op_data->object_state->osd_set) + { + // Check is the same as in submit_primary_del_subops() + if (pg.scheme == POOL_SCHEME_REPLICATED + ? !contains_osd(pg.cur_set.data(), pg.pg_size, chunk.osd_num) + : (chunk.osd_num != pg.cur_set[chunk.role])) + { + pg.copies_to_delete_after_sync.push_back((obj_ver_osd_t){ + .osd_num = chunk.osd_num, + .oid = { + .inode = op_data->oid.inode, + .stripe = op_data->oid.stripe | (pg.scheme == POOL_SCHEME_REPLICATED ? 0 : chunk.role), + }, + .version = op_data->fact_ver, + }); + copies_to_delete_after_sync_count++; + } + } + if (pg.copies_to_delete_after_sync.size() && !(pg.state & PG_HAS_MISPLACED)) + { + // PG can't be active+clean until extra copies aren't removed, so mark it as PG_HAS_MISPLACED + pg.state |= PG_HAS_MISPLACED; + //this->pg_state_dirty.insert({ .pool_id = pg.pool_id, .pg_num = pg.pg_num }); + } + } // We must forget the unclean state of the object before deleting it // so the next reads don't accidentally read a deleted version // And it should be done at the same time as the removal of the version override @@ -309,6 +341,7 @@ resume_12: } resume_6: resume_7: + op_data->n_subops = 0; if (!remember_unstable_write(cur_op, pg, pg.cur_loc_set, 6)) { return; @@ -344,48 +377,21 @@ resume_7: ); recovery_stat[recovery_type].usec += usec; } - // Any kind of a non-clean object can have extra chunks, because we don't record objects - // as degraded & misplaced or incomplete & misplaced at the same time. So try to remove extra chunks - if (immediate_commit != IMMEDIATE_ALL) - { - // We can't remove extra chunks yet if fsyncs are explicit, because - // new copies may not be committed to stable storage yet - // We can only remove extra chunks after a successful SYNC for this PG - for (auto & chunk: op_data->object_state->osd_set) - { - // Check is the same as in submit_primary_del_subops() - if (pg.scheme == POOL_SCHEME_REPLICATED - ? !contains_osd(pg.cur_set.data(), pg.pg_size, chunk.osd_num) - : (chunk.osd_num != pg.cur_set[chunk.role])) - { - pg.copies_to_delete_after_sync.push_back((obj_ver_osd_t){ - .osd_num = chunk.osd_num, - .oid = { - .inode = op_data->oid.inode, - .stripe = op_data->oid.stripe | (pg.scheme == POOL_SCHEME_REPLICATED ? 0 : chunk.role), - }, - .version = op_data->fact_ver, - }); - copies_to_delete_after_sync_count++; - } - } - deref_object_state(pg, &op_data->object_state, true); - } - else + if (immediate_commit == IMMEDIATE_ALL) { submit_primary_del_subops(cur_op, pg.cur_set.data(), pg.pg_size, op_data->object_state->osd_set); - deref_object_state(pg, &op_data->object_state, true); - if (op_data->n_subops > 0) - { + } + deref_object_state(pg, &op_data->object_state, true); + if (op_data->n_subops > 0) + { resume_8: - op_data->st = 8; - return; + op_data->st = 8; + return; resume_9: - if (op_data->errors > 0) - { - pg_cancel_write_queue(pg, cur_op, op_data->oid, op_data->errcode); - return; - } + if (op_data->errors > 0) + { + pg_cancel_write_queue(pg, cur_op, op_data->oid, op_data->errcode); + return; } } } diff --git a/src/osd/osd_secondary.cpp b/src/osd/osd_secondary.cpp index 135cd293..93dda04b 100644 --- a/src/osd/osd_secondary.cpp +++ b/src/osd/osd_secondary.cpp @@ -79,6 +79,32 @@ 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) +{ + if (!enable_pg_locks) + { + return true; + } + 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()) + { + return false; + } + auto ppg = (pool_pg_num_t){ .pool_id = pool_id, .pg_num = map_to_pg(oid, pool_cfg_it->second.pg_stripe_size) }; + auto pg_it = pgs.find(ppg); + if (pg_it != pgs.end() && pg_it->second.state != PG_OFFLINE) + { + return false; + } + if (pg_it->second.disable_pg_locks) + { + return true; + } + auto lock_it = pg_locks.find(ppg); + return lock_it != pg_locks.end() && lock_it->second.primary_osd == primary_osd; +} + void osd_t::exec_secondary_real(osd_op_t *cur_op) { if (cur_op->req.hdr.opcode == OSD_OP_SEC_LIST && @@ -89,23 +115,15 @@ void osd_t::exec_secondary_real(osd_op_t *cur_op) } if (cur_op->req.hdr.opcode == OSD_OP_SEC_READ_BMP) { - int n = cur_op->req.sec_read_bmp.len / sizeof(obj_ver_id); - if (n > 0) - { - obj_ver_id *ov = (obj_ver_id*)cur_op->buf; - void *reply_buf = malloc_or_die(n * (8 + clean_entry_bitmap_size)); - void *cur_buf = reply_buf; - for (int i = 0; i < n; i++) - { - bs->read_bitmap(ov[i].oid, ov[i].version, (uint8_t*)cur_buf + sizeof(uint64_t), (uint64_t*)cur_buf); - cur_buf = (uint8_t*)cur_buf + (8 + clean_entry_bitmap_size); - } - free(cur_op->buf); - cur_op->buf = reply_buf; - } - finish_op(cur_op, n * (8 + clean_entry_bitmap_size)); + exec_sec_read_bmp(cur_op); return; } + else if (cur_op->req.hdr.opcode == OSD_OP_SEC_LOCK) + { + exec_sec_lock(cur_op); + return; + } + auto cl = msgr.clients.at(cur_op->peer_fd); cur_op->bs_op = new blockstore_op_t(); cur_op->bs_op->callback = [this, cur_op](blockstore_op_t* bs_op) { secondary_op_callback(cur_op); }; cur_op->bs_op->opcode = (cur_op->req.hdr.opcode == OSD_OP_SEC_READ ? BS_OP_READ @@ -121,6 +139,13 @@ void osd_t::exec_secondary_real(osd_op_t *cur_op) cur_op->req.hdr.opcode == OSD_OP_SEC_WRITE || cur_op->req.hdr.opcode == OSD_OP_SEC_WRITE_STABLE) { + if (!(cur_op->req.sec_rw.flags & OSD_OP_IGNORE_PG_LOCK) && + !sec_check_pg_lock(cl->osd_num, cur_op->req.sec_rw.oid)) + { + cur_op->bs_op->retval = -EPIPE; + secondary_op_callback(cur_op); + return; + } if (cur_op->req.hdr.opcode == OSD_OP_SEC_READ) { // Allocate memory for the read operation @@ -143,6 +168,13 @@ void osd_t::exec_secondary_real(osd_op_t *cur_op) } else if (cur_op->req.hdr.opcode == OSD_OP_SEC_DELETE) { + if (!(cur_op->req.sec_del.flags & OSD_OP_IGNORE_PG_LOCK) && + !sec_check_pg_lock(cl->osd_num, cur_op->req.sec_del.oid)) + { + cur_op->bs_op->retval = -EPIPE; + secondary_op_callback(cur_op); + return; + } cur_op->bs_op->oid = cur_op->req.sec_del.oid; cur_op->bs_op->version = cur_op->req.sec_del.version; #ifdef OSD_STUB @@ -157,6 +189,18 @@ void osd_t::exec_secondary_real(osd_op_t *cur_op) #ifdef OSD_STUB cur_op->bs_op->retval = 0; #endif + if (enable_pg_locks && !(cur_op->req.sec_stab.flags & OSD_OP_IGNORE_PG_LOCK)) + { + for (int i = 0; i < cur_op->bs_op->len; i++) + { + if (!sec_check_pg_lock(cl->osd_num, ((obj_ver_id*)cur_op->buf)[i].oid)) + { + cur_op->bs_op->retval = -EPIPE; + secondary_op_callback(cur_op); + return; + } + } + } } else if (cur_op->req.hdr.opcode == OSD_OP_SEC_LIST) { @@ -192,15 +236,96 @@ void osd_t::exec_secondary_real(osd_op_t *cur_op) #endif } +void osd_t::exec_sec_read_bmp(osd_op_t *cur_op) +{ + auto cl = msgr.clients.at(cur_op->peer_fd); + int n = cur_op->req.sec_read_bmp.len / sizeof(obj_ver_id); + if (n > 0) + { + obj_ver_id *ov = (obj_ver_id*)cur_op->buf; + void *reply_buf = malloc_or_die(n * (8 + clean_entry_bitmap_size)); + void *cur_buf = reply_buf; + for (int i = 0; i < n; i++) + { + if (!sec_check_pg_lock(cl->osd_num, ov[i].oid) && + !(cur_op->req.sec_read_bmp.flags & OSD_OP_IGNORE_PG_LOCK)) + { + free(reply_buf); + cur_op->bs_op->retval = -EPIPE; + secondary_op_callback(cur_op); + return; + } + bs->read_bitmap(ov[i].oid, ov[i].version, (uint8_t*)cur_buf + sizeof(uint64_t), (uint64_t*)cur_buf); + cur_buf = (uint8_t*)cur_buf + (8 + clean_entry_bitmap_size); + } + free(cur_op->buf); + cur_op->buf = reply_buf; + } + finish_op(cur_op, n * (8 + clean_entry_bitmap_size)); +} + +// Lock/Unlock PG +void osd_t::exec_sec_lock(osd_op_t *cur_op) +{ + cur_op->reply.sec_lock.cur_primary = 0; + auto cl = msgr.clients.at(cur_op->peer_fd); + if (!cl->osd_num || + cur_op->req.sec_lock.flags != OSD_SEC_LOCK_PG && + cur_op->req.sec_lock.flags != OSD_SEC_UNLOCK_PG || + cur_op->req.sec_lock.pool_id > ((uint64_t)1<req.sec_lock.pg_num || + cur_op->req.sec_lock.pg_num > UINT32_MAX) + { + finish_op(cur_op, -EINVAL); + 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() || + pool_cfg_it->second.real_pg_count < cur_op->req.sec_lock.pg_num) + { + finish_op(cur_op, -ENOENT); + return; + } + auto lock_it = pg_locks.find(ppg); + if (cur_op->req.sec_lock.flags == OSD_SEC_LOCK_PG) + { + if (lock_it != pg_locks.end() && lock_it->second.primary_osd != cl->osd_num) + { + cur_op->reply.sec_lock.cur_primary = lock_it->second.primary_osd; + finish_op(cur_op, -EBUSY); + return; + } + auto primary_pg_it = pgs.find(ppg); + if (primary_pg_it != pgs.end() && primary_pg_it->second.state != PG_OFFLINE) + { + cur_op->reply.sec_lock.cur_primary = this->osd_num; + finish_op(cur_op, -EBUSY); + return; + } + pg_locks[ppg] = (osd_pg_lock_t){ + .primary_osd = cl->osd_num, + .state = cur_op->req.sec_lock.pg_state, + }; + } + else if (lock_it != pg_locks.end() && lock_it->second.primary_osd == cl->osd_num) + { + pg_locks.erase(lock_it); + } + finish_op(cur_op, 0); +} + void osd_t::exec_show_config(osd_op_t *cur_op) { std::string json_err; json11::Json req_json = cur_op->req.show_conf.json_len > 0 ? json11::Json::parse(std::string((char *)cur_op->buf), json_err) : json11::Json(); + auto peer_osd_num = req_json["osd_num"].uint64_value(); + auto cl = msgr.clients.at(cur_op->peer_fd); + cl->osd_num = peer_osd_num; if (req_json["features"]["check_sequencing"].bool_value()) { - auto cl = msgr.clients.at(cur_op->peer_fd); cl->check_sequencing = true; cl->read_op_id = cur_op->req.hdr.id + 1; } @@ -216,6 +341,7 @@ void osd_t::exec_show_config(osd_op_t *cur_op) { "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 }, + { "features", json11::Json::object{ { "pg_locks", true } } }, }; #ifdef WITH_RDMA if (msgr.is_rdma_enabled()) @@ -228,7 +354,7 @@ void osd_t::exec_show_config(osd_op_t *cur_op) bool ok = msgr.connect_rdma(cur_op->peer_fd, req_json["connect_rdma"].string_value(), req_json["rdma_max_msg"].uint64_value()); if (ok) { - auto rc = msgr.clients.at(cur_op->peer_fd)->rdma_conn; + auto rc = cl->rdma_conn; wire_config["rdma_address"] = rc->addr.to_string(); wire_config["rdma_max_msg"] = rc->max_msg; }