// Copyright (c) Vitaliy Filippov, 2019+ // License: VNPL-1.1 or GNU GPL-2.0+ (see README.md for details) #include #include "assert.h" #include "pg_states.h" #include "cluster_client.h" #define LIST_PG_INIT 0 #define LIST_PG_WAIT_ACTIVE 1 #define LIST_PG_WAIT_CONNECT 2 #define LIST_PG_SENT 3 #define LIST_PG_DONE 4 struct inode_list_t; struct inode_list_pg_t; struct inode_list_osd_t { inode_list_pg_t *pg = NULL; osd_num_t osd_num = 0; }; struct inode_list_pg_t { inode_list_t *lst = NULL; int errcode = 0; pg_num_t pg_num = 0; osd_num_t cur_primary = 0; int state = 0; int inflight_ops = 0; timespec wait_until; std::vector list_osds; bool has_unstable = false; std::set objects; std::vector inactive_osds; }; struct inode_list_t { cluster_client_t *cli = NULL; pool_id_t pool_id = 0; inode_t inode = 0; int max_parallel_pgs = 16; int inflight_pgs = 0; std::map inflight_per_osd; int done_pgs = 0; int onstack = 0; std::vector pgs; pg_num_t real_pg_count = 0; std::function&& objects, pg_num_t pg_num, std::vector && inactive_osds, int errcode, int status)> callback; }; inode_list_t* cluster_client_t::list_inode_start(inode_t inode, int max_parallel_pgs, std::function&& objects, pg_num_t pg_num, std::vector && inactive_osds, int errcode, int status)> callback) { init_msgr(); pool_id_t pool_id = INODE_POOL(inode); if (!pool_id || st_cli.pool_config.find(pool_id) == st_cli.pool_config.end()) { if (log_level > 0) { fprintf(stderr, "Pool %u does not exist\n", pool_id); } return NULL; } inode_list_t *lst = new inode_list_t(); lst->cli = this; lst->pool_id = pool_id; lst->inode = inode; lst->callback = callback; lst->max_parallel_pgs = max_parallel_pgs <= 0 ? 16 : max_parallel_pgs; lists.push_back(lst); if (!continue_listing(lst)) { return NULL; } return lst; } int cluster_client_t::list_pg_count(inode_list_t *lst) { return lst->pgs.size() - lst->done_pgs; } void cluster_client_t::list_inode_next(inode_list_t *lst) { continue_listing(lst); } bool cluster_client_t::continue_listing(inode_list_t *lst) { if (lst->onstack > 0) { return true; } lst->onstack++; if (restart_listing(lst)) { for (int i = 0; i < lst->pgs.size() && lst->inflight_pgs < lst->max_parallel_pgs; i++) { retry_start_pg_listing(lst->pgs[i]); } } if (check_finish_listing(lst)) { // Do not change lst->onstack because it's already freed return false; } lst->onstack--; return true; } bool cluster_client_t::restart_listing(inode_list_t* lst) { auto pool_it = st_cli.pool_config.find(lst->pool_id); // We want listing to be consistent. To achieve it we should: // 1) retry listing of each PG if its state changes // 2) abort listing if PG count changes during listing // 3) ideally, only talk to the primary OSD - this will be done separately // So first we add all PGs without checking their state if (pool_it == st_cli.pool_config.end() || lst->real_pg_count != pool_it->second.real_pg_count) { for (auto pg: lst->pgs) { if (pg->inflight_ops > 0) { // Wait until all in-progress listings complete or fail return false; } } for (auto pg: lst->pgs) { delete pg; } if (log_level > 0 && lst->real_pg_count) { fprintf(stderr, "PG count in pool %u changed during listing\n", lst->pool_id); } lst->pgs.clear(); if (pool_it == st_cli.pool_config.end()) { // Unknown pool lst->callback(lst, std::set(), 0, std::vector(), -EINVAL, INODE_LIST_DONE); return false; } else if (lst->done_pgs) { // PG count changed during listing, it should fail lst->callback(lst, std::set(), 0, std::vector(), -EAGAIN, INODE_LIST_DONE); return false; } else { lst->real_pg_count = pool_it->second.real_pg_count; for (pg_num_t pg_num = 1; pg_num <= lst->real_pg_count; pg_num++) { inode_list_pg_t *pg = new inode_list_pg_t(); pg->lst = lst; pg->pg_num = pg_num; lst->pgs.push_back(pg); } } } return true; } void cluster_client_t::retry_start_pg_listing(inode_list_pg_t *pg) { if (pg->state == LIST_PG_SENT || pg->state == LIST_PG_DONE) { return; } int new_st = start_pg_listing(pg); if (new_st == LIST_PG_SENT) { pg->state = LIST_PG_SENT; return; } if (new_st == LIST_PG_WAIT_ACTIVE && pg->state != LIST_PG_WAIT_ACTIVE || new_st == LIST_PG_WAIT_CONNECT && pg->state != LIST_PG_WAIT_CONNECT) { int sec = (new_st == LIST_PG_WAIT_ACTIVE ? wait_up_timeout : peer_connect_timeout); if (sec) { pg->state = new_st; clock_gettime(CLOCK_REALTIME, &pg->wait_until); pg->wait_until.tv_sec += sec; if (new_st == LIST_PG_WAIT_ACTIVE) { if (log_level > 1) fprintf(stderr, "Waiting for PG %u/%u to become active for %d seconds\n", pg->lst->pool_id, pg->pg_num, wait_up_timeout); } else { if (log_level > 2) fprintf(stderr, "Waiting for connection to PG %u/%u OSDs for %d seconds\n", pg->lst->pool_id, pg->pg_num, peer_connect_timeout); } if (!list_retry_time.tv_sec || list_retry_time.tv_sec > pg->wait_until.tv_sec || list_retry_time.tv_sec == pg->wait_until.tv_sec && list_retry_time.tv_nsec > pg->wait_until.tv_nsec) { list_retry_time = pg->wait_until; if (list_retry_timeout_id >= 0) { tfd->clear_timer(list_retry_timeout_id); } list_retry_timeout_id = tfd->set_timer(sec*1000, false, [this](int timer_id) { continue_lists(); }); } return; } } assert(pg->state == LIST_PG_WAIT_ACTIVE || pg->state == LIST_PG_WAIT_CONNECT); // Check if the timeout expired timespec tv; clock_gettime(CLOCK_REALTIME, &tv); if (tv.tv_sec > pg->wait_until.tv_sec || tv.tv_sec == pg->wait_until.tv_sec && tv.tv_nsec >= pg->wait_until.tv_nsec) { fprintf(stderr, "Failed to wait for PG %u/%u to become %s, skipping listing\n", pg->lst->pool_id, pg->pg_num, pg->state == LIST_PG_WAIT_ACTIVE ? "active" : "connected"); pg->errcode = -EPIPE; pg->list_osds.clear(); pg->objects.clear(); finish_list_pg(pg); } } int cluster_client_t::start_pg_listing(inode_list_pg_t *pg) { auto & pool_cfg = st_cli.pool_config[pg->lst->pool_id]; auto pg_it = pool_cfg.pg_config.find(pg->pg_num); assert(pg->lst->real_pg_count == pool_cfg.real_pg_count); if (pg_it == pool_cfg.pg_config.end() || pg_it->second.pause || !pg_it->second.cur_primary || !(pg_it->second.cur_state & PG_ACTIVE)) { // PG is (temporarily?) unavailable return LIST_PG_WAIT_ACTIVE; } pg->inactive_osds.clear(); std::set all_peers; if (pg_it->second.cur_state != PG_ACTIVE) { // Not clean for (osd_num_t pg_osd: pg_it->second.target_set) all_peers.insert(pg_osd); for (osd_num_t pg_osd: pg_it->second.all_peers) all_peers.insert(pg_osd); for (auto & hist_item: pg_it->second.target_history) for (auto pg_osd: hist_item) all_peers.insert(pg_osd); // Remove zero OSD number all_peers.erase(0); // Remove unconnectable peers except cur_primary for (auto peer_it = all_peers.begin(); peer_it != all_peers.end(); ) { if (*peer_it != pg_it->second.cur_primary && st_cli.peer_states[*peer_it].is_null()) { pg->inactive_osds.push_back(*peer_it); printf("OSD is inactive: %lu\n", *peer_it); all_peers.erase(peer_it++); } else peer_it++; } } else { // Clean all_peers.insert(pg_it->second.cur_primary); } // Check that we're connected to all PG OSDs bool conn = true; for (osd_num_t peer_osd: all_peers) { if (msgr.osd_peer_fds.find(peer_osd) == msgr.osd_peer_fds.end()) { // Initiate connection if (st_cli.peer_states[peer_osd].is_null()) { return LIST_PG_WAIT_ACTIVE; } msgr.connect_peer(peer_osd, st_cli.peer_states[peer_osd]); conn = false; } } if (!conn) { return LIST_PG_WAIT_CONNECT; } // Send all listings at once as the simplest way to guarantee that we connect // to the exact same OSDs that are listed in PG state pg->errcode = 0; pg->list_osds.clear(); pg->has_unstable = false; pg->objects.clear(); pg->cur_primary = pg_it->second.cur_primary; for (osd_num_t peer_osd: all_peers) { pg->list_osds.push_back((inode_list_osd_t){ .pg = pg, .osd_num = peer_osd, }); } for (auto & list_osd: pg->list_osds) { send_list(&list_osd); } return LIST_PG_SENT; } void cluster_client_t::send_list(inode_list_osd_t *cur_list) { if (!cur_list->pg->inflight_ops) cur_list->pg->lst->inflight_pgs++; cur_list->pg->inflight_ops++; auto & pool_cfg = st_cli.pool_config[cur_list->pg->lst->pool_id]; osd_op_t *op = new osd_op_t(); op->op_type = OSD_OP_OUT; // Already checked that it exists above, but anyway op->peer_fd = msgr.osd_peer_fds.at(cur_list->osd_num); op->req = (osd_any_op_t){ .sec_list = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, .id = next_op_id(), .opcode = OSD_OP_SEC_LIST, }, .list_pg = cur_list->pg->pg_num, .pg_count = (pg_num_t)pool_cfg.real_pg_count, .pg_stripe_size = pool_cfg.pg_stripe_size, .min_inode = cur_list->pg->lst->inode, .max_inode = cur_list->pg->lst->inode, }, }; op->callback = [this, cur_list](osd_op_t *op) { if (op->reply.hdr.retval < 0) { fprintf(stderr, "Failed to get PG %u/%u object list from OSD %ju (retval=%jd), skipping\n", cur_list->pg->lst->pool_id, cur_list->pg->pg_num, cur_list->osd_num, op->reply.hdr.retval); cur_list->pg->errcode = op->reply.hdr.retval; } else { if (op->reply.sec_list.stable_count < op->reply.hdr.retval) { // Unstable objects, if present, mean that someone still writes into the inode. Warn the user about it. cur_list->pg->has_unstable = true; fprintf( stderr, "[PG %u/%u] Inode still has %ju unstable object versions out of total %ju - is it still open?\n", cur_list->pg->lst->pool_id, cur_list->pg->pg_num, op->reply.hdr.retval - op->reply.sec_list.stable_count, op->reply.hdr.retval ); } if (log_level > 0) { fprintf( stderr, "[PG %u/%u] Got inode object list from OSD %ju: %jd object versions\n", cur_list->pg->lst->pool_id, cur_list->pg->pg_num, cur_list->osd_num, op->reply.hdr.retval ); } for (uint64_t i = 0; i < op->reply.hdr.retval; i++) { object_id oid = ((obj_ver_id*)op->buf)[i].oid; oid.stripe = oid.stripe & ~STRIPE_MASK; cur_list->pg->objects.insert(oid); } } delete op; cur_list->pg->inflight_ops--; if (!cur_list->pg->inflight_ops) cur_list->pg->lst->inflight_pgs--; // FIXME: Retry listing after ms on EPIPE finish_list_pg(cur_list->pg); continue_listing(cur_list->pg->lst); }; msgr.outbox_push(op); } void cluster_client_t::finish_list_pg(inode_list_pg_t *pg) { auto lst = pg->lst; if (pg->inflight_ops == 0) { lst->done_pgs++; pg->state = LIST_PG_DONE; int status = 0; if (lst->done_pgs >= lst->pgs.size()) { status |= INODE_LIST_DONE; } if (pg->has_unstable) { status |= INODE_LIST_HAS_UNSTABLE; } lst->callback(lst, std::move(pg->objects), pg->pg_num, std::move(pg->inactive_osds), pg->errcode, status); } } void cluster_client_t::continue_lists() { for (int i = lists.size()-1; i >= 0; i--) { continue_listing(lists[i]); } } bool cluster_client_t::check_finish_listing(inode_list_t *lst) { if (lst->done_pgs >= lst->pgs.size()) { for (auto pg: lst->pgs) { delete pg; } lst->pgs.clear(); for (int i = 0; i < lists.size(); i++) { if (lists[i] == lst) { lists.erase(lists.begin()+i, lists.begin()+i+1); break; } } delete lst; return true; } return false; }