From 2ab0ae3bc96b61b08e5e6c3a472b50ebb33e8f40 Mon Sep 17 00:00:00 2001 From: Vitaliy Filippov Date: Wed, 16 Apr 2025 13:54:05 +0300 Subject: [PATCH] Check operation sequencing and stop clients when it breaks --- src/client/cluster_client.cpp | 7 ------- src/client/cluster_client.h | 1 - src/client/cluster_client_list.cpp | 1 - src/client/messenger.cpp | 28 +++++++++++++++++----------- src/client/messenger.h | 4 +++- src/client/msgr_receive.cpp | 12 ++++++++++++ src/client/msgr_send.cpp | 1 + src/cmd/cli_describe.cpp | 1 - src/cmd/cli_fix.cpp | 3 --- src/osd/osd_flush.cpp | 1 - src/osd/osd_peering.cpp | 1 - src/osd/osd_primary_chain.cpp | 1 - src/osd/osd_primary_subops.cpp | 5 ----- src/osd/osd_scrub.cpp | 1 - src/osd/osd_secondary.cpp | 6 ++++++ src/test/mock/messenger.cpp | 4 +++- 16 files changed, 42 insertions(+), 35 deletions(-) diff --git a/src/client/cluster_client.cpp b/src/client/cluster_client.cpp index e89dbd9f..a83a0758 100644 --- a/src/client/cluster_client.cpp +++ b/src/client/cluster_client.cpp @@ -1244,7 +1244,6 @@ int cluster_client_t::try_send(cluster_op_t *op, int i) .req = { .rw = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = next_op_id(), .opcode = op->opcode == OSD_OP_READ_BITMAP || op->opcode == OSD_OP_READ_CHAIN_BITMAP ? OSD_OP_READ : op->opcode, }, .inode = op->cur_inode, @@ -1353,7 +1352,6 @@ void cluster_client_t::send_sync(cluster_op_t *op, cluster_op_part_t *part) .req = { .hdr = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = next_op_id(), .opcode = OSD_OP_SYNC, }, }, @@ -1498,8 +1496,3 @@ void cluster_client_t::copy_part_bitmap(cluster_op_t *op, cluster_op_part_t *par part_len--; } } - -uint64_t cluster_client_t::next_op_id() -{ - return msgr.next_subop_id++; -} diff --git a/src/client/cluster_client.h b/src/client/cluster_client.h index d4f2357d..cb7ef496 100644 --- a/src/client/cluster_client.h +++ b/src/client/cluster_client.h @@ -152,7 +152,6 @@ public: //inline uint32_t get_bs_bitmap_granularity() { return st_cli.global_bitmap_granularity; } //inline uint64_t get_bs_block_size() { return st_cli.global_block_size; } - uint64_t next_op_id(); #ifndef __MOCK__ protected: diff --git a/src/client/cluster_client_list.cpp b/src/client/cluster_client_list.cpp index c7898db6..3d4cd6f9 100644 --- a/src/client/cluster_client_list.cpp +++ b/src/client/cluster_client_list.cpp @@ -342,7 +342,6 @@ void cluster_client_t::send_list(inode_list_osd_t *cur_list) .sec_list = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = next_op_id(), .opcode = OSD_OP_SEC_LIST, }, .list_pg = cur_list->pg->pg_num, diff --git a/src/client/messenger.cpp b/src/client/messenger.cpp index 8afc9b47..78fd78cb 100644 --- a/src/client/messenger.cpp +++ b/src/client/messenger.cpp @@ -217,7 +217,6 @@ void osd_messenger_t::init() op->req = (osd_any_op_t){ .hdr = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = this->next_subop_id++, .opcode = OSD_OP_PING, }, }; @@ -629,11 +628,17 @@ void osd_messenger_t::check_peer_config(osd_client_t *cl) .show_conf = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = this->next_subop_id++, .opcode = OSD_OP_SHOW_CONFIG, }, }, }; + json11::Json::object payload; + if (osd_num) + { + // Inform that we're OSD + payload["osd_num"] = osd_num; + } + payload["check_sequencing"] = true; #ifdef WITH_RDMA if (!use_rdmacm && rdma_contexts.size()) { @@ -649,19 +654,20 @@ void osd_messenger_t::check_peer_config(osd_client_t *cl) cl->rdma_conn = msgr_rdma_connection_t::create(selected_ctx, rdma_max_send, rdma_max_recv, rdma_max_sge, rdma_max_msg); if (cl->rdma_conn) { - json11::Json payload = json11::Json::object { - { "connect_rdma", cl->rdma_conn->addr.to_string() }, - { "rdma_max_msg", cl->rdma_conn->max_msg }, - }; - std::string payload_str = payload.dump(); - op->req.show_conf.json_len = payload_str.size(); - op->buf = malloc_or_die(payload_str.size()); - op->iov.push_back(op->buf, payload_str.size()); - memcpy(op->buf, payload_str.c_str(), payload_str.size()); + payload["connect_rdma"] = cl->rdma_conn->addr.to_string(); + payload["rdma_max_msg"] = cl->rdma_conn->max_msg; } } } #endif + if (payload.size()) + { + std::string payload_str = json11::Json(payload).dump(); + op->req.show_conf.json_len = payload_str.size(); + op->buf = malloc_or_die(payload_str.size()); + op->iov.push_back(op->buf, payload_str.size()); + memcpy(op->buf, payload_str.c_str(), payload_str.size()); + } op->callback = [this, cl](osd_op_t *op) { std::string json_err; diff --git a/src/client/messenger.h b/src/client/messenger.h index 3cf93707..1754a666 100644 --- a/src/client/messenger.h +++ b/src/client/messenger.h @@ -75,12 +75,15 @@ struct osd_client_t int read_remaining = 0; int read_state = 0; osd_op_buf_list_t recv_list; + uint64_t read_op_id = 1; + bool check_sequencing = false; // Incoming operations std::vector received_ops; // Outbound operations std::map sent_ops; + uint64_t send_op_id = 0; // PGs dirtied by this client's primary-writes std::set dirty_pgs; @@ -211,7 +214,6 @@ public: bool has_sendmsg_zc = false; // osd_num_t is only for logging and asserts osd_num_t osd_num; - uint64_t next_subop_id = 1; std::map clients; std::map wanted_peers; std::map osd_peer_fds; diff --git a/src/client/msgr_receive.cpp b/src/client/msgr_receive.cpp index 011f4611..51b431d9 100644 --- a/src/client/msgr_receive.cpp +++ b/src/client/msgr_receive.cpp @@ -223,7 +223,19 @@ bool osd_messenger_t::handle_finished_read(osd_client_t *cl) if (cl->read_op->req.hdr.magic == SECONDARY_OSD_REPLY_MAGIC) return handle_reply_hdr(cl); else if (cl->read_op->req.hdr.magic == SECONDARY_OSD_OP_MAGIC) + { + if (cl->check_sequencing) + { + if (cl->read_op->req.hdr.id != cl->read_op_id) + { + fprintf(stderr, "Warning: operation sequencing is broken on client %d, stopping client\n", cl->peer_fd); + stop_client(cl->peer_fd); + return false; + } + cl->read_op_id++; + } handle_op_hdr(cl); + } else { fprintf(stderr, "Received garbage: magic=%jx id=%ju opcode=%jx from %d\n", cl->read_op->req.hdr.magic, cl->read_op->req.hdr.id, cl->read_op->req.hdr.opcode, cl->peer_fd); diff --git a/src/client/msgr_send.cpp b/src/client/msgr_send.cpp index c4e573b3..c9d0fdec 100644 --- a/src/client/msgr_send.cpp +++ b/src/client/msgr_send.cpp @@ -14,6 +14,7 @@ void osd_messenger_t::outbox_push(osd_op_t *cur_op) if (cur_op->op_type == OSD_OP_OUT) { clock_gettime(CLOCK_REALTIME, &cur_op->tv_begin); + cur_op->req.hdr.id = ++cl->send_op_id; } else { diff --git a/src/cmd/cli_describe.cpp b/src/cmd/cli_describe.cpp index a15c94fb..7f489d47 100644 --- a/src/cmd/cli_describe.cpp +++ b/src/cmd/cli_describe.cpp @@ -147,7 +147,6 @@ struct cli_describe_t .describe = (osd_op_describe_t){ .header = (osd_op_header_t){ .magic = SECONDARY_OSD_OP_MAGIC, - .id = parent->cli->next_op_id(), .opcode = OSD_OP_DESCRIBE, }, .object_state = object_state, diff --git a/src/cmd/cli_fix.cpp b/src/cmd/cli_fix.cpp index 9016fcfb..6f57dc4c 100644 --- a/src/cmd/cli_fix.cpp +++ b/src/cmd/cli_fix.cpp @@ -159,7 +159,6 @@ struct cli_fix_t .describe = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = parent->cli->next_op_id(), .opcode = OSD_OP_DESCRIBE, }, .min_inode = obj.inode, @@ -194,7 +193,6 @@ struct cli_fix_t .sec_del = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = parent->cli->next_op_id(), .opcode = OSD_OP_SEC_DELETE, }, .oid = { @@ -242,7 +240,6 @@ struct cli_fix_t .rw = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = parent->cli->next_op_id(), .opcode = OSD_OP_SCRUB, }, .inode = obj.inode, diff --git a/src/osd/osd_flush.cpp b/src/osd/osd_flush.cpp index b697c644..a663188e 100644 --- a/src/osd/osd_flush.cpp +++ b/src/osd/osd_flush.cpp @@ -209,7 +209,6 @@ bool osd_t::submit_flush_op(pool_id_t pool_id, pg_num_t pg_num, pg_flush_batch_t .sec_stab = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = msgr.next_subop_id++, .opcode = (uint64_t)(rollback ? OSD_OP_SEC_ROLLBACK : OSD_OP_SEC_STABILIZE), }, .len = count * sizeof(obj_ver_id), diff --git a/src/osd/osd_peering.cpp b/src/osd/osd_peering.cpp index 67db9196..e70fc6f0 100644 --- a/src/osd/osd_peering.cpp +++ b/src/osd/osd_peering.cpp @@ -391,7 +391,6 @@ void osd_t::submit_list_subop(osd_num_t role_osd, pg_peering_state_t *ps) .sec_list = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = msgr.next_subop_id++, .opcode = OSD_OP_SEC_LIST, }, .list_pg = ps->pg_num, diff --git a/src/osd/osd_primary_chain.cpp b/src/osd/osd_primary_chain.cpp index 94375daa..62d7fc9e 100644 --- a/src/osd/osd_primary_chain.cpp +++ b/src/osd/osd_primary_chain.cpp @@ -266,7 +266,6 @@ int osd_t::submit_bitmap_subops(osd_op_t *cur_op, pg_t & pg) .sec_read_bmp = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = msgr.next_subop_id++, .opcode = OSD_OP_SEC_READ_BMP, }, .len = sizeof(obj_ver_id)*(i+1-prev), diff --git a/src/osd/osd_primary_subops.cpp b/src/osd/osd_primary_subops.cpp index 97a8f70d..3549d078 100644 --- a/src/osd/osd_primary_subops.cpp +++ b/src/osd/osd_primary_subops.cpp @@ -233,7 +233,6 @@ void osd_t::submit_primary_subop(osd_op_t *cur_op, osd_op_t *subop, subop->req.sec_rw = (osd_op_sec_rw_t){ .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = msgr.next_subop_id++, .opcode = (uint64_t)(wr ? (cur_op->op_data->pg->scheme == POOL_SCHEME_REPLICATED ? OSD_OP_SEC_WRITE_STABLE : OSD_OP_SEC_WRITE) : OSD_OP_SEC_READ), }, .oid = { @@ -594,7 +593,6 @@ void osd_t::submit_primary_del_batch(osd_op_t *cur_op, obj_ver_osd_t *chunks_to_ subops[i].req = (osd_any_op_t){ .sec_del = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = msgr.next_subop_id++, .opcode = OSD_OP_SEC_DELETE, }, .oid = chunk.oid, @@ -654,7 +652,6 @@ int osd_t::submit_primary_sync_subops(osd_op_t *cur_op) subops[i].req = (osd_any_op_t){ .sec_sync = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = msgr.next_subop_id++, .opcode = OSD_OP_SEC_SYNC, }, .flags = cur_op->peer_fd == SELF_FD && cur_op->req.hdr.opcode != OSD_OP_SCRUB ? OSD_OP_RECOVERY_RELATED : 0, @@ -713,7 +710,6 @@ void osd_t::submit_primary_stab_subops(osd_op_t *cur_op) subops[i].req = (osd_any_op_t){ .sec_stab = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = msgr.next_subop_id++, .opcode = OSD_OP_SEC_STABILIZE, }, .len = (uint64_t)(stab_osd.len * sizeof(obj_ver_id)), @@ -807,7 +803,6 @@ void osd_t::submit_primary_rollback_subops(osd_op_t *cur_op, const uint64_t* osd subop->req = (osd_any_op_t){ .sec_stab = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = msgr.next_subop_id++, .opcode = OSD_OP_SEC_ROLLBACK, }, .len = sizeof(obj_ver_id), diff --git a/src/osd/osd_scrub.cpp b/src/osd/osd_scrub.cpp index bee5616a..b4480068 100644 --- a/src/osd/osd_scrub.cpp +++ b/src/osd/osd_scrub.cpp @@ -65,7 +65,6 @@ void osd_t::scrub_list(pool_pg_num_t pg_id, osd_num_t role_osd, object_id min_oi .sec_list = { .header = { .magic = SECONDARY_OSD_OP_MAGIC, - .id = msgr.next_subop_id++, .opcode = OSD_OP_SEC_LIST, }, .list_pg = pg_num, diff --git a/src/osd/osd_secondary.cpp b/src/osd/osd_secondary.cpp index dec606f5..a5d1a43f 100644 --- a/src/osd/osd_secondary.cpp +++ b/src/osd/osd_secondary.cpp @@ -198,6 +198,12 @@ void osd_t::exec_show_config(osd_op_t *cur_op) 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(); + if (req_json["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; + } // Expose sensitive configuration values so peers can check them json11::Json::object wire_config = json11::Json::object { { "osd_num", osd_num }, diff --git a/src/test/mock/messenger.cpp b/src/test/mock/messenger.cpp index 2d5004b5..8443f702 100644 --- a/src/test/mock/messenger.cpp +++ b/src/test/mock/messenger.cpp @@ -21,7 +21,9 @@ osd_messenger_t::~osd_messenger_t() void osd_messenger_t::outbox_push(osd_op_t *cur_op) { - clients[cur_op->peer_fd]->sent_ops[cur_op->req.hdr.id] = cur_op; + auto cl = clients.at(cur_op->peer_fd); + cur_op->req.hdr.id = ++cl->send_op_id; + cl->sent_ops[cur_op->req.hdr.id] = cur_op; } void osd_messenger_t::parse_config(const json11::Json & config)