Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions include/pingcap/Config.h
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@

#include <kvproto/kvrpcpb.pb.h>

#include <cstdint>
#include <fstream>
#include <streambuf>
#include <string>
Expand All @@ -23,6 +24,8 @@ struct ClusterConfig
std::string cert_path;
std::string key_path;
::kvrpcpb::APIVersion api_version = ::kvrpcpb::APIVersion::V1;
::kvrpcpb::RequestOrigin request_origin = ::kvrpcpb::RequestOriginUnknown;
uint32_t default_txn_protocol_version = ::kvrpcpb::TXN_VER_SUPPORT_INCOMPATIBLE_ERROR_HANDLING;

ClusterConfig() = default;

Expand Down
40 changes: 39 additions & 1 deletion include/pingcap/Exception.h
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
#pragma once

#include <Poco/Exception.h>
#include <kvproto/errorpb.pb.h>

#include <exception>
#include <string>
Expand Down Expand Up @@ -29,7 +30,9 @@ enum ErrorCodes : int
KeyspaceNotEnabled = 18,
InternalError = 19,
GRPCNotImplemented = 20,
UnknownError = 21
UnknownError = 21,
IncompatibleRequest = 22,
UndeterminedResult = 23,
};

class Exception : public Poco::Exception
Expand All @@ -55,6 +58,41 @@ class Exception : public Poco::Exception
bool empty() const { return code() == 0 && message().empty(); }
};

class ErrIncompatibleRequest : public Exception
{
public:
explicit ErrIncompatibleRequest(const ::errorpb::IncompatibleRequest & error)
: Exception(error.message(), ErrorCodes::IncompatibleRequest)
, error_(error)
{}

ErrIncompatibleRequest(const std::string & message, const ::errorpb::IncompatibleRequest & error)
: Exception(message, ErrorCodes::IncompatibleRequest)
, error_(error)
{}

const ::errorpb::IncompatibleRequest & error() const { return error_; }

ErrIncompatibleRequest * clone() const override { return new ErrIncompatibleRequest(*this); }
void rethrow() const override { throw *this; }

private:
::errorpb::IncompatibleRequest error_;
};

inline bool isTerminalTransactionError(const Exception & exception)
{
return exception.code() == ErrorCodes::IncompatibleRequest || exception.code() == ErrorCodes::UndeterminedResult;
}

inline void rethrowTerminalRegionError(const ::errorpb::Error & error)
{
if (error.has_undetermined_result())
throw Exception(error.undetermined_result().message(), UndeterminedResult);
if (error.has_incompatible_request())
throw ErrIncompatibleRequest(error.incompatible_request());
}

inline std::string getCurrentExceptionMsg(const std::string & prefix_msg)
{
std::string msg = prefix_msg;
Expand Down
7 changes: 7 additions & 0 deletions include/pingcap/coprocessor/Client.h
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
#include <atomic>
#include <condition_variable>
#include <cstdint>
#include <memory>
#include <mutex>
#include <thread>

Expand Down Expand Up @@ -113,6 +114,9 @@ class ResponseIter
std::shared_ptr<::coprocessor::Response> resp;
bool same_zone{true};
Exception error;
// Preserve the dynamic exception and, for ErrIncompatibleRequest, its
// complete structured protobuf across the asynchronous queue boundary.
std::shared_ptr<Exception> detailed_error;
bool finished{false};

Result() = default;
Expand All @@ -121,6 +125,7 @@ class ResponseIter
{}
explicit Result(const Exception & err)
: error(err)
, detailed_error(err.clone())
{}
explicit Result(bool finished_)
: finished(finished_)
Expand All @@ -131,6 +136,8 @@ class ResponseIter
{}

const std::string & data() const { return resp->data(); }

const Exception * exception() const { return detailed_error.get(); }
};

ResponseIter(std::unique_ptr<common::IMPMCQueue<Result>> && queue_,
Expand Down
18 changes: 17 additions & 1 deletion include/pingcap/kv/Cluster.h
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,8 @@ struct Cluster
LockResolverPtr lock_resolver;

::kvrpcpb::APIVersion api_version = ::kvrpcpb::APIVersion::V1;
const ::kvrpcpb::RequestOrigin request_origin;
const uint32_t default_txn_protocol_version;

std::unique_ptr<pingcap::common::FixedThreadPool> thread_pool;
std::unique_ptr<common::MPPProber> mpp_prober;
Expand All @@ -39,19 +41,23 @@ struct Cluster
, rpc_client(std::make_unique<RpcClient>(pd_client, ClusterConfig{}))
, oracle(std::make_unique<pd::Oracle>(pd_client, std::chrono::milliseconds(oracle_update_interval)))
, lock_resolver(std::make_unique<LockResolver>(this))
, request_origin(ClusterConfig{}.request_origin)
, default_txn_protocol_version(ClusterConfig{}.default_txn_protocol_version)
, thread_pool(std::make_unique<pingcap::common::FixedThreadPool>(mock_cluster_background_workers))
, mpp_prober(std::make_unique<common::MPPProber>(this))
{
startBackgroundTasks();
}

Cluster(const std::vector<std::string> & pd_addrs, const ClusterConfig & config)
: pd_client(std::make_shared<pd::CodecClient>(pd_addrs, config))
: pd_client(std::make_shared<pd::CodecClient>(pd_addrs, validateCompatibilityConfig(config)))
, region_cache(std::make_unique<RegionCache>(pd_client, config))
, rpc_client(std::make_unique<RpcClient>(pd_client, config))
, oracle(std::make_unique<pd::Oracle>(pd_client, std::chrono::milliseconds(oracle_update_interval)))
, lock_resolver(std::make_unique<LockResolver>(this))
, api_version(config.api_version)
, request_origin(config.request_origin)
, default_txn_protocol_version(config.default_txn_protocol_version)
, thread_pool(std::make_unique<pingcap::common::FixedThreadPool>(cluster_background_workers))
, mpp_prober(std::make_unique<common::MPPProber>(this))
{
Expand All @@ -60,6 +66,8 @@ struct Cluster

void update(const std::vector<std::string> & pd_addrs, const ClusterConfig & config) const
{
if (config.request_origin != request_origin || config.default_txn_protocol_version != default_txn_protocol_version)
throw Exception("request origin and transaction protocol version are immutable after Cluster creation", LogicalError);
pd_client->update(pd_addrs, config);
rpc_client->update(config);
}
Expand All @@ -81,6 +89,14 @@ struct Cluster
void splitRegion(const std::string & split_key);

void startBackgroundTasks();

private:
static const ClusterConfig & validateCompatibilityConfig(const ClusterConfig & config)
{
if (config.default_txn_protocol_version > ::kvrpcpb::TXN_VER_SUPPORT_INCOMPATIBLE_ERROR_HANDLING)
throw Exception("default transaction protocol version must be legacy or incompatible-error-handling", LogicalError);
return config;
}
};

struct MinCommitTSPushed
Expand Down
9 changes: 8 additions & 1 deletion include/pingcap/kv/RegionCache.h
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@

#include <atomic>
#include <chrono>
#include <cstdint>
#include <map>
#include <unordered_map>

Expand All @@ -30,14 +31,20 @@ struct Store
std::map<std::string, std::string> labels;
StoreType store_type;
::metapb::StoreState state;
bool has_txn_protocol_version_range;
uint32_t txn_protocol_version_min;
uint32_t txn_protocol_version_max;

Store(uint64_t id_, const std::string & addr_, const std::string & peer_addr_, const std::map<std::string, std::string> & labels_, StoreType store_type_, const ::metapb::StoreState state_)
Store(uint64_t id_, const std::string & addr_, const std::string & peer_addr_, const std::map<std::string, std::string> & labels_, StoreType store_type_, const ::metapb::StoreState state_, bool has_txn_protocol_version_range_ = false, uint32_t txn_protocol_version_min_ = 0, uint32_t txn_protocol_version_max_ = 0)
: id(id_)
, addr(addr_)
, peer_addr(peer_addr_)
, labels(labels_)
, store_type(store_type_)
, state(state_)
, has_txn_protocol_version_range(has_txn_protocol_version_range_)
, txn_protocol_version_min(txn_protocol_version_min_)
, txn_protocol_version_max(txn_protocol_version_max_)
{}
};

Expand Down
104 changes: 100 additions & 4 deletions include/pingcap/kv/RegionClient.h
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
#include <pingcap/kv/Cluster.h>
#include <pingcap/kv/RegionCache.h>
#include <pingcap/kv/Rpc.h>
#include <pingcap/kv/internal/txn_protocol.h>

namespace pingcap
{
Expand Down Expand Up @@ -53,9 +54,15 @@ struct RegionClient
{
throw Exception("should setup proper label_filter for tiflash");
}
bool compatibility_resend_used = false;
RPCContextPtr compatibility_resend_ctx;
uint32_t compatibility_resend_version = 0;
for (;;)
{
RPCContextPtr ctx = cluster->region_cache->getRPCContext(bo, region_id, store_type, /*load_balance=*/true, tiflash_label_filter, store_id_blocklist, prefer_store_id);
const bool is_compatibility_resend = compatibility_resend_ctx != nullptr;
RPCContextPtr ctx = is_compatibility_resend
? compatibility_resend_ctx
: cluster->region_cache->getRPCContext(bo, region_id, store_type, /*load_balance=*/true, tiflash_label_filter, store_id_blocklist, prefer_store_id);
if (ctx == nullptr)
{
// If the region is not found in cache, it must be out
Expand All @@ -64,8 +71,19 @@ struct RegionClient
auto s = store_id_blocklist != nullptr ? ", store_filter_size=" + std::to_string(store_id_blocklist->size()) + "." : std::string(".");
throw Exception("Region epoch not match after retries: Region " + region_id.toString() + " not in region cache" + s, RegionEpochNotMatch);
}
auto selection = internal::selectTxnProtocolVersion(
req, cluster->default_txn_protocol_version, ctx->store.txn_protocol_version_min, ctx->store.txn_protocol_version_max);
if (is_compatibility_resend)
{
selection.selected = compatibility_resend_version;
selection.allowed = !selection.protected_request || selection.required <= selection.selected;
compatibility_resend_ctx.reset();
}
if (!selection.allowed)
throw localIncompatibleRequest(ctx, selection);

RpcCall<T> rpc(cluster->rpc_client, ctx->addr);
rpc.setRequestCtx(req, ctx, cluster->api_version);
rpc.setRequestCtx(req, ctx, cluster->api_version, cluster->request_origin, selection.selected);

grpc::ClientContext context;
rpc.setClientContext(context, timeout, meta_data);
Expand All @@ -87,6 +105,29 @@ struct RegionClient
if (resp->has_region_error())
{
log->warning("region_id " + region_id.toString() + " find error: " + resp->region_error().DebugString());
const auto & error = resp->region_error();
if (error.has_undetermined_result())
throw Exception(error.undetermined_result().message(), UndeterminedResult);
if (error.has_incompatible_request())
{
const auto & incompatible = error.incompatible_request();
if (!compatibility_resend_used && !is_compatibility_resend
&& internal::isValidUpperBoundRejection(incompatible, selection.selected))
{
auto updated = internal::selectTxnProtocolVersion(
req,
cluster->default_txn_protocol_version,
incompatible.min_compatible_txn_protocol_version(),
incompatible.max_compatible_txn_protocol_version());
if (updated.allowed && updated.selected != selection.selected)
{
compatibility_resend_used = true;
compatibility_resend_ctx = ctx;
compatibility_resend_version = updated.selected;
continue;
}
}
}
onRegionError(bo, ctx, resp->region_error());
continue;
}
Expand Down Expand Up @@ -153,9 +194,15 @@ struct RegionClient
{
throw Exception("should setup proper label_filter for tiflash");
}
bool compatibility_resend_used = false;
RPCContextPtr compatibility_resend_ctx;
uint32_t compatibility_resend_version = 0;
for (;;)
{
RPCContextPtr ctx = cluster->region_cache->getRPCContext(bo, region_id, store_type, /*load_balance=*/true, tiflash_label_filter, store_id_blocklist, prefer_store_id);
const bool is_compatibility_resend = compatibility_resend_ctx != nullptr;
RPCContextPtr ctx = is_compatibility_resend
? compatibility_resend_ctx
: cluster->region_cache->getRPCContext(bo, region_id, store_type, /*load_balance=*/true, tiflash_label_filter, store_id_blocklist, prefer_store_id);
if (ctx == nullptr)
{
// If the region is not found in cache, it must be out
Expand All @@ -164,9 +211,20 @@ struct RegionClient
throw Exception("Region epoch not match after retries: Region " + region_id.toString() + " not in region cache.", RegionEpochNotMatch);
}

auto selection = internal::selectTxnProtocolVersion(
req, cluster->default_txn_protocol_version, ctx->store.txn_protocol_version_min, ctx->store.txn_protocol_version_max);
if (is_compatibility_resend)
{
selection.selected = compatibility_resend_version;
selection.allowed = !selection.protected_request || selection.required <= selection.selected;
compatibility_resend_ctx.reset();
}
if (!selection.allowed)
throw localIncompatibleRequest(ctx, selection);

auto stream_reader = std::make_unique<StreamReader<RESP>>();
RpcCall<T> rpc(cluster->rpc_client, ctx->addr);
rpc.setRequestCtx(req, ctx, cluster->api_version);
rpc.setRequestCtx(req, ctx, cluster->api_version, cluster->request_origin, selection.selected);
rpc.setClientContext(stream_reader->context, timeout, meta_data);

stream_reader->reader = rpc.call(&stream_reader->context, req);
Expand All @@ -175,6 +233,29 @@ struct RegionClient
if (stream_reader->first_resp.has_region_error())
{
log->warning("region_id " + region_id.toString() + " find error: " + stream_reader->first_resp.region_error().message());
const auto & error = stream_reader->first_resp.region_error();
if (error.has_undetermined_result())
throw Exception(error.undetermined_result().message(), UndeterminedResult);
if (error.has_incompatible_request())
{
const auto & incompatible = error.incompatible_request();
if (!compatibility_resend_used && !is_compatibility_resend
&& internal::isValidUpperBoundRejection(incompatible, selection.selected))
{
auto updated = internal::selectTxnProtocolVersion(
req,
cluster->default_txn_protocol_version,
incompatible.min_compatible_txn_protocol_version(),
incompatible.max_compatible_txn_protocol_version());
if (updated.allowed && updated.selected != selection.selected)
{
compatibility_resend_used = true;
compatibility_resend_ctx = ctx;
compatibility_resend_version = updated.selected;
continue;
}
}
}
onRegionError(bo, ctx, stream_reader->first_resp.region_error());
continue;
}
Expand Down Expand Up @@ -207,6 +288,21 @@ struct RegionClient
}

protected:
static ErrIncompatibleRequest localIncompatibleRequest(const RPCContextPtr & ctx, const internal::TxnProtocolSelection & selection)
{
::errorpb::IncompatibleRequest error;
error.set_reason(::errorpb::IncompatibleRequestReasonUnknown);
error.set_min_compatible_txn_protocol_version(ctx->store.txn_protocol_version_min);
error.set_max_compatible_txn_protocol_version(ctx->store.txn_protocol_version_max);
error.set_provided_txn_protocol_version(selection.selected);
const auto message = "transaction protocol is incompatible with store " + std::to_string(ctx->store.id) + " range ["
+ std::to_string(ctx->store.txn_protocol_version_min) + "," + std::to_string(ctx->store.txn_protocol_version_max)
+ "], process ceiling " + std::to_string(selection.process_ceiling) + ", candidate " + std::to_string(selection.selected)
+ ", required " + std::to_string(selection.required);
error.set_message(message);
return ErrIncompatibleRequest(message, error);
}

void onRegionError(Backoffer & bo, RPCContextPtr rpc_ctx, const errorpb::Error & err) const;

// Normally, it happens when machine down or network partition between tidb and kv or process crash.
Expand Down
Loading
Loading