diff --git a/src/client/osd_ops.h b/src/client/osd_ops.h index 9fd7be9a..e67ed8b6 100644 --- a/src/client/osd_ops.h +++ b/src/client/osd_ops.h @@ -51,6 +51,8 @@ #define LOC_CORRUPTED 2 #define LOC_INCONSISTENT 4 +#define OSD_LIST_PRIMARY 1 + // common request and reply headers struct __attribute__((__packed__)) osd_op_header_t { @@ -196,6 +198,9 @@ struct __attribute__((__packed__)) osd_op_sec_list_t uint64_t min_stripe, max_stripe; // max stable object count uint32_t stable_limit; + // flags - OSD_LIST_PRIMARY or 0 + // for OSD_LIST_PRIMARY, only a single-PG listing is allowed + uint64_t flags; }; struct __attribute__((__packed__)) osd_reply_sec_list_t @@ -204,6 +209,8 @@ struct __attribute__((__packed__)) osd_reply_sec_list_t // stable object version count. header.retval = total object version count // FIXME: maybe change to the number of bytes in the reply... uint64_t stable_count; + // flags - OSD_LIST_PRIMARY or 0 + uint64_t flags; }; // read or write to the primary OSD (must be within individual stripe) diff --git a/src/osd/osd.h b/src/osd/osd.h index 43091076..33eb7694 100644 --- a/src/osd/osd.h +++ b/src/osd/osd.h @@ -301,6 +301,7 @@ class osd_t void continue_primary_read(osd_op_t *cur_op); void continue_primary_scrub(osd_op_t *cur_op); void continue_primary_describe(osd_op_t *cur_op); + void continue_primary_list(osd_op_t *cur_op); void continue_primary_write(osd_op_t *cur_op); void cancel_primary_write(osd_op_t *cur_op); void continue_primary_sync(osd_op_t *cur_op); diff --git a/src/osd/osd_primary_describe.cpp b/src/osd/osd_primary_describe.cpp index 1fce43ab..80ec92c7 100644 --- a/src/osd/osd_primary_describe.cpp +++ b/src/osd/osd_primary_describe.cpp @@ -138,3 +138,108 @@ void osd_t::continue_primary_describe(osd_op_t *cur_op) cur_op->iov.push_back(res.items, res.size * sizeof(osd_reply_describe_item_t)); finish_op(cur_op, res.size); } + +static void add_primary_list(btree::btree_map & list, osd_op_sec_list_t & req, std::set & oids) +{ + auto begin_it = list.begin(); + auto end_it = list.end(); + if (req.min_inode) + begin_it = list.lower_bound((object_id){ .inode = req.min_inode, .stripe = req.min_stripe }); + if (req.max_inode) + end_it = list.upper_bound((object_id){ .inode = req.max_inode, .stripe = (req.max_stripe ? req.max_stripe : UINT64_MAX) }); + for (auto list_it = begin_it; list_it != end_it; list_it++) + oids.insert(list_it->first); +} + +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)); + // 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() || + (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)) + { + finish_op(cur_op, -EINVAL); + return; + } + auto pg_it = pgs.find({ .pool_id = INODE_POOL(cur_op->req.sec_list.min_inode), .pg_num = cur_op->req.sec_list.list_pg }); + if (pg_it == pgs.end()) + { + // Not primary + finish_op(cur_op, -EPIPE); + return; + } + cur_op->bs_op = new blockstore_op_t(); + cur_op->bs_op->opcode = BS_OP_LIST; + cur_op->bs_op->pg_alignment = pool_cfg_it->second.pg_stripe_size; + cur_op->bs_op->pg_count = pool_cfg_it->second.real_pg_count; + cur_op->bs_op->pg_number = cur_op->req.sec_list.list_pg - 1; + cur_op->bs_op->min_oid.inode = cur_op->req.sec_list.min_inode; + cur_op->bs_op->min_oid.stripe = cur_op->req.sec_list.min_stripe; + cur_op->bs_op->max_oid.inode = cur_op->req.sec_list.max_inode; + if (cur_op->req.sec_list.max_inode && cur_op->req.sec_list.max_stripe != UINT64_MAX) + { + cur_op->bs_op->max_oid.stripe = cur_op->req.sec_list.max_stripe + ? cur_op->req.sec_list.max_stripe : UINT64_MAX; + } + cur_op->bs_op->callback = [this, cur_op](blockstore_op_t* bs_op) + { + if (bs_op->retval < 0) + { + auto retval = bs_op->retval; + delete bs_op; + cur_op->bs_op = NULL; + finish_op(cur_op, retval); + return; + } + // Move into a set + std::set oids; + obj_ver_id *rbuf = (obj_ver_id*)bs_op->buf; + uint64_t total_count = bs_op->retval; + for (uint64_t i = 0; i < total_count; i++) + { + oids.insert(rbuf[i].oid); + } + if (bs_op->buf) + { + free(bs_op->buf); + bs_op->buf = NULL; + } + delete bs_op; + cur_op->bs_op = NULL; + // Check if still primary + auto pg_it = pgs.find({ .pool_id = INODE_POOL(cur_op->req.sec_list.min_inode), .pg_num = cur_op->req.sec_list.list_pg }); + if (pg_it == pgs.end()) + { + finish_op(cur_op, -EPIPE); + return; + } + // Add unclean objects which may be not present on the primary OSD + auto & pg = pg_it->second; + add_primary_list(pg.inconsistent_objects, cur_op->req.sec_list, oids); + add_primary_list(pg.incomplete_objects, cur_op->req.sec_list, oids); + add_primary_list(pg.degraded_objects, cur_op->req.sec_list, oids); + add_primary_list(pg.misplaced_objects, cur_op->req.sec_list, oids); + // Generate the result + if (oids.size()) + { + rbuf = (obj_ver_id*)malloc_or_die(sizeof(obj_ver_id) * oids.size()); + uint64_t i = 0; + for (auto oid: oids) + { + rbuf[i++] = (obj_ver_id){ .oid = oid }; + } + cur_op->buf = rbuf; + cur_op->iov.push_back(rbuf, sizeof(obj_ver_id) * oids.size()); + } + cur_op->reply.sec_list.stable_count = oids.size(); + cur_op->reply.sec_list.flags = OSD_LIST_PRIMARY; + finish_op(cur_op, oids.size()); + }; + bs->enqueue_op(cur_op->bs_op); +} diff --git a/src/osd/osd_secondary.cpp b/src/osd/osd_secondary.cpp index 61c26d87..dec606f5 100644 --- a/src/osd/osd_secondary.cpp +++ b/src/osd/osd_secondary.cpp @@ -81,6 +81,12 @@ void osd_t::exec_secondary(osd_op_t *op) void osd_t::exec_secondary_real(osd_op_t *cur_op) { + if (cur_op->req.hdr.opcode == OSD_OP_SEC_LIST && + (cur_op->req.sec_list.flags & OSD_LIST_PRIMARY)) + { + continue_primary_list(cur_op); + return; + } if (cur_op->req.hdr.opcode == OSD_OP_SEC_READ_BMP) { int n = cur_op->req.sec_read_bmp.len / sizeof(obj_ver_id);