| 667 | |
| 668 | template<typename Storage> |
| 669 | nuraft::ptr<nuraft::buffer> KeeperStateMachine<Storage>::commit(const uint64_t log_idx, nuraft::buffer & data) |
| 670 | { |
| 671 | const UInt64 start_time_us = ZooKeeperOpentelemetrySpans::now(); |
| 672 | |
| 673 | auto request_for_session = parseRequest(data, true); |
| 674 | if (!request_for_session->zxid) |
| 675 | request_for_session->zxid = log_idx; |
| 676 | |
| 677 | request_for_session->log_idx = log_idx; |
| 678 | |
| 679 | if (!keeper_context->localLogsPreprocessed() && !preprocess(*request_for_session, /*lock_mutex=*/ true)) |
| 680 | return nullptr; |
| 681 | |
| 682 | const auto maybe_log_opentelemetry_span = [&](OpenTelemetry::SpanStatus status, const std::string & error_message) |
| 683 | { |
| 684 | request_for_session->request->spans.maybeInitialize( |
| 685 | KeeperSpan::Commit, |
| 686 | request_for_session->request->tracing_context.get(), |
| 687 | start_time_us); |
| 688 | |
| 689 | request_for_session->request->spans.maybeFinalize( |
| 690 | KeeperSpan::Commit, |
| 691 | [&] |
| 692 | { |
| 693 | return std::vector<OpenTelemetry::SpanAttribute>{ |
| 694 | {"keeper.operation", Coordination::opNumToString(request_for_session->request->getOpNum())}, |
| 695 | {"keeper.session_id", request_for_session->session_id}, |
| 696 | {"keeper.xid", request_for_session->request->xid}, |
| 697 | {"raft.log_idx", log_idx}, |
| 698 | }; |
| 699 | }, |
| 700 | status, |
| 701 | error_message); |
| 702 | }; |
| 703 | |
| 704 | try |
| 705 | { |
| 706 | const auto op_num = request_for_session->request->getOpNum(); |
| 707 | if (op_num == Coordination::OpNum::SessionID) |
| 708 | { |
| 709 | const Coordination::ZooKeeperSessionIDRequest & session_id_request |
| 710 | = dynamic_cast<const Coordination::ZooKeeperSessionIDRequest &>(*request_for_session->request); |
| 711 | int64_t session_id = 0; |
| 712 | std::shared_ptr<Coordination::ZooKeeperSessionIDResponse> response = std::dynamic_pointer_cast<Coordination::ZooKeeperSessionIDResponse>(session_id_request.makeResponse()); |
| 713 | KeeperResponseForSession response_for_session; |
| 714 | response_for_session.session_id = -1; |
| 715 | response_for_session.response = response; |
| 716 | response_for_session.request = request_for_session->request; |
| 717 | |
| 718 | KEEPER_STORAGE_LOCK_EXCLUSIVE(lock); |
| 719 | session_id = storage->getSessionID(session_id_request.session_timeout_ms); |
| 720 | LOG_DEBUG(log, "Session ID response {} with timeout {}", session_id, session_id_request.session_timeout_ms); |
| 721 | response->session_id = session_id; |
| 722 | if (response_callback) |
| 723 | response_callback(std::move(response_for_session)); |
| 724 | } |
| 725 | else |
| 726 | { |
nothing calls this directly
no test coverage detected