| 365 | */ |
| 366 | |
| 367 | AsyncMessenger::AsyncMessenger(CephContext *cct, entity_name_t name, |
| 368 | const std::string &type, std::string mname, uint64_t _nonce) |
| 369 | : SimplePolicyMessenger(cct, name), |
| 370 | dispatch_queue(cct, this, mname), |
| 371 | nonce(_nonce) |
| 372 | { |
| 373 | std::string transport_type = "posix"; |
| 374 | if (type.find("rdma") != std::string::npos) |
| 375 | transport_type = "rdma"; |
| 376 | else if (type.find("dpdk") != std::string::npos) |
| 377 | transport_type = "dpdk"; |
| 378 | else if (type.find("smc") != std::string::npos) |
| 379 | transport_type = "smc"; |
| 380 | |
| 381 | auto single = &cct->lookup_or_create_singleton_object<StackSingleton>( |
| 382 | "AsyncMessenger::NetworkStack::" + transport_type, true, cct); |
| 383 | single->ready(transport_type); |
| 384 | stack = single->stack.get(); |
| 385 | stack->start(); |
| 386 | local_worker = stack->get_worker(); |
| 387 | local_connection = ceph::make_ref<AsyncConnection>(cct, this, &dispatch_queue, |
| 388 | local_worker, true, true); |
| 389 | init_local_connection(); |
| 390 | reap_handler = new C_handle_reap(this); |
| 391 | unsigned processor_num = 1; |
| 392 | if (stack->support_local_listen_table()) |
| 393 | processor_num = stack->get_num_worker(); |
| 394 | for (unsigned i = 0; i < processor_num; ++i) |
| 395 | processors.push_back(new Processor(this, stack->get_worker(i), cct)); |
| 396 | |
| 397 | cct->modify_msgr_hook( |
| 398 | [&]() { |
| 399 | auto hook = new AsyncMessengerSocketHook(*this, mname); |
| 400 | const int asok_ret = cct->get_admin_socket()->register_command( |
| 401 | AsyncMessengerSocketHook::COMMAND, hook, "dump messenger status"); |
| 402 | if (asok_ret != 0) { |
| 403 | ldout(cct, 0) << __func__ << " messenger asok command \"" |
| 404 | << AsyncMessengerSocketHook::COMMAND << "\" failed with" |
| 405 | << asok_ret << dendl; |
| 406 | } |
| 407 | return hook; |
| 408 | }, |
| 409 | [&](AdminSocketHook* ptr) { |
| 410 | if (auto hook = dynamic_cast<AsyncMessengerSocketHook*>(ptr)) { |
| 411 | // Name collisions may occur. A common case is RGW running |
| 412 | // multiple librados connections - each with a "radosclient" |
| 413 | // messenger. |
| 414 | std::string msgr_key(mname); |
| 415 | if (mname == "radosclient") { // librados |
| 416 | msgr_key.append("-"); |
| 417 | msgr_key.append(std::to_string(_nonce)); |
| 418 | } |
| 419 | bool added = hook->add_messenger(msgr_key, *this); |
| 420 | if (!added) { |
| 421 | msgr_key.append("-"); |
| 422 | msgr_key.append(std::to_string(_nonce)); |
| 423 | } |
| 424 | added = hook->add_messenger(msgr_key, *this); |
nothing calls this directly
no test coverage detected