MCPcopy Create free account
hub / github.com/ClickHouse/ClickHouse / commit

Method commit

src/Coordination/KeeperStateMachine.cpp:669–777  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

667
668template<typename Storage>
669nuraft::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 {

Callers

nothing calls this directly

Calls 15

preprocessFunction · 0.85
opNumToStringFunction · 0.85
assertDigestFunction · 0.85
localLogsPreprocessedMethod · 0.80
maybeInitializeMethod · 0.80
maybeFinalizeMethod · 0.80
makeResponseMethod · 0.80
digestEnabledMethod · 0.80
setLastCommitIndexMethod · 0.80
nowFunction · 0.50
incrementFunction · 0.50

Tested by

no test coverage detected