Skip to content

Commit de385f7

Browse files
LiuRuoyu01chuandew
authored andcommitted
[fix][sdk]Fixup get commit ts opportunity
1 parent e3a4105 commit de385f7

5 files changed

Lines changed: 136 additions & 65 deletions

File tree

src/sdk/transaction/txn_impl.cc

Lines changed: 18 additions & 43 deletions
Original file line numberDiff line numberDiff line change
@@ -568,52 +568,27 @@ Status TxnImpl::PreCommit2PC() {
568568
return Status::OK();
569569
}
570570

571-
Status TxnImpl::ProcessTxnCommitResponse(const TxnCommitResponse* response, bool is_primary) {
572-
std::string pk = buffer_->GetPrimaryKey();
573-
DINGO_LOG(DEBUG) << fmt::format("[sdk.txn.{}] commit response, pk({}) response({}).", ID(), pk,
574-
response->ShortDebugString());
575-
576-
if (!response->has_txn_result()) {
577-
return Status::OK();
578-
}
579-
580-
const auto& txn_result = response->txn_result();
581-
if (txn_result.has_locked()) {
582-
const auto& lock_info = txn_result.locked();
583-
DINGO_LOG(FATAL) << fmt::format("[sdk.txn.{}] commit lock conflict, is_primary({}) pk({}) response({}).", ID(),
584-
is_primary, StringToHex(pk), response->ShortDebugString());
585-
586-
} else if (txn_result.has_txn_not_found()) {
587-
DINGO_LOG(FATAL) << fmt::format("[sdk.txn.{}] commit not found, is_primary({}) pk({}) response({}).", ID(),
588-
is_primary, StringToHex(pk), response->ShortDebugString());
589-
590-
} else if (txn_result.has_write_conflict()) {
591-
if (!is_primary) {
592-
DINGO_LOG(FATAL) << fmt::format("[sdk.txn.{}] commit write conlict, pk({}) response({}).", ID(), StringToHex(pk),
593-
txn_result.write_conflict().ShortDebugString());
594-
}
595-
return Status::TxnWriteConflict("txn write conflict");
596-
597-
} else if (txn_result.has_commit_ts_expired()) {
598-
DINGO_LOG(WARNING) << fmt::format("[sdk.txn.{}] commit ts expired, is_primary({}) pk({}) response({}).", ID(),
599-
is_primary, StringToHex(pk), txn_result.commit_ts_expired().ShortDebugString());
600-
if (is_primary) {
601-
int64_t new_commit_ts;
602-
auto status = stub_.GetTsoProvider()->GenTs(2, new_commit_ts);
603-
commit_ts_.store(new_commit_ts);
604-
if (!status.IsOK()) return status;
605-
return Status::TxnCommitTsExpired("txn commit ts expired");
606-
}
607-
}
608-
609-
return Status::OK();
610-
}
611-
612571
Status TxnImpl::CommitPrimaryKey() {
613572
std::vector<std::string> keys = {buffer_->GetPrimaryKey()};
573+
Status status;
574+
int64_t retry_count = 0;
575+
do {
576+
TxnCommitTask task(stub_, keys, shared_from_this(), true);
577+
status = task.Run();
578+
if (status.IsTxnCommitTsExpired()) {
579+
int64_t commit_ts;
580+
Status s = stub_.GetTsoProvider()->GenTs(2, commit_ts);
581+
if (!s.ok()) {
582+
DINGO_LOG(ERROR) << fmt::format("[sdk.txn.{}] commit primary key regen ts fail, status({}).", ID(),
583+
s.ToString());
584+
return s;
585+
}
586+
commit_ts_.store(commit_ts);
587+
}
588+
retry_count++;
589+
} while (status.IsTxnCommitTsExpired() && retry_count < FLAGS_txn_op_max_retry);
614590

615-
TxnCommitTask task(stub_, keys, shared_from_this(), true);
616-
return task.Run();
591+
return status;
617592
}
618593

619594
Status TxnImpl::CommitOrdinaryKey() {

src/sdk/transaction/txn_impl.h

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -127,12 +127,11 @@ class TxnImpl : public std::enable_shared_from_this<TxnImpl> {
127127

128128
Status Rollback();
129129

130-
Status ProcessTxnCommitResponse(const TxnCommitResponse* response, bool is_primary);
131-
132130
bool IsOnePc() const { return is_one_pc_; }
133131

134132
int64_t GetStartTs() const { return start_ts_.load(); }
135133
int64_t GetCommitTs() const { return commit_ts_.load(); }
134+
std::string GetPrimaryKey() const { return buffer_->GetPrimaryKey(); }
136135
TransactionOptions GetOptions() const { return options_; }
137136

138137
bool CheckFinished() const {

src/sdk/transaction/txn_task/txn_commit_task.cc

Lines changed: 46 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -17,11 +17,13 @@
1717
#include <fmt/format.h>
1818
#include <glog/logging.h>
1919

20+
#include <cstdint>
2021
#include <utility>
2122

2223
#include "common/logging.h"
2324
#include "dingosdk/status.h"
2425
#include "sdk/common/common.h"
26+
#include "sdk/common/helper.h"
2527
#include "sdk/rpc/store_rpc_controller.h"
2628
#include "sdk/transaction/txn_common.h"
2729
#include "sdk/utils/rw_lock.h"
@@ -55,7 +57,6 @@ void TxnCommitTask::DoAsync() {
5557
{
5658
WriteLockGuard guard(rw_lock_);
5759
next_batch = next_keys_;
58-
need_retry_ = false;
5960
status_ = Status::OK();
6061
}
6162

@@ -136,37 +137,27 @@ void TxnCommitTask::TxnCommitRpcCallback(const Status& status, TxnCommitRpc* rpc
136137
DINGO_LOG(DEBUG) << fmt::format("[sdk.txn.{}] rpc: {} request: {} response: {}", txn_impl_->ID(), rpc->Method(),
137138
rpc->Request()->ShortDebugString(), rpc->Response()->ShortDebugString());
138139
Status s;
139-
bool need_retry = false;
140140
const auto* response = rpc->Response();
141141
if (!status.ok()) {
142142
DINGO_LOG(WARNING) << fmt::format("[sdk.txn.{}] rpc: {} send to region: {} fail: {}", txn_impl_->ID(),
143143
rpc->Method(), rpc->Request()->context().region_id(), status.ToString());
144144

145145
s = status;
146146
} else {
147-
s = txn_impl_->ProcessTxnCommitResponse(response, is_primary_);
147+
s = ProcessTxnCommitResponse(response, is_primary_);
148148
if (!s.ok()) {
149149
DINGO_LOG(WARNING) << fmt::format("[sdk.txn.{}] commit fail, region({}) response({}) status({}).",
150150
txn_impl_->ID(), rpc->Request()->context().region_id(),
151151
response->ShortDebugString(), s.ToString());
152-
if (s.IsTxnCommitTsExpired() && is_primary_) {
153-
need_retry = true;
154-
s = Status::OK();
155-
}
156152
}
157153
}
158154

159155
{
160156
WriteLockGuard guard(rw_lock_);
161157
if (s.ok()) {
162-
if (!need_retry) {
163-
for (const auto& key : rpc->Request()->keys()) {
164-
next_keys_.erase(key);
165-
}
166-
} else {
167-
need_retry_ = true;
158+
for (const auto& key : rpc->Request()->keys()) {
159+
next_keys_.erase(key);
168160
}
169-
170161
} else {
171162
if (status_.ok()) {
172163
// only return first fail status
@@ -181,16 +172,52 @@ void TxnCommitTask::TxnCommitRpcCallback(const Status& status, TxnCommitRpc* rpc
181172
{
182173
ReadLockGuard guard(rw_lock_);
183174
tmp = status_;
184-
tmp_need_retry = need_retry_;
185-
}
186-
if (tmp.ok() && tmp_need_retry) {
187-
DoAsyncRetry();
188-
return;
189175
}
190176
DoAsyncDone(tmp);
191177
}
192178
}
193179

180+
Status TxnCommitTask::ProcessTxnCommitResponse(const TxnCommitResponse* response, bool is_primary) {
181+
std::string pk = txn_impl_->GetPrimaryKey();
182+
int64_t txn_id = txn_impl_->ID();
183+
DINGO_LOG(DEBUG) << fmt::format("[sdk.txn.{}] commit response, pk({}) response({}).", txn_id, pk,
184+
response->ShortDebugString());
185+
186+
if (!response->has_txn_result()) {
187+
return Status::OK();
188+
}
189+
190+
const auto& txn_result = response->txn_result();
191+
if (txn_result.has_locked()) {
192+
const auto& lock_info = txn_result.locked();
193+
DINGO_LOG(FATAL) << fmt::format("[sdk.txn.{}] commit lock conflict, is_primary({}) pk({}) response({}).", txn_id,
194+
is_primary, StringToHex(pk), response->ShortDebugString());
195+
196+
} else if (txn_result.has_txn_not_found()) {
197+
DINGO_LOG(FATAL) << fmt::format("[sdk.txn.{}] commit not found, is_primary({}) pk({}) response({}).", txn_id,
198+
is_primary, StringToHex(pk), response->ShortDebugString());
199+
200+
} else if (txn_result.has_write_conflict()) {
201+
if (!is_primary) {
202+
DINGO_LOG(FATAL) << fmt::format("[sdk.txn.{}] commit write conlict, pk({}) response({}).", txn_id,
203+
StringToHex(pk), txn_result.write_conflict().ShortDebugString());
204+
}
205+
return Status::TxnWriteConflict("txn write conflict");
206+
207+
} else if (txn_result.has_commit_ts_expired()) {
208+
DINGO_LOG(WARNING) << fmt::format("[sdk.txn.{}] commit ts expired, is_primary({}) pk({}) response({}).", txn_id,
209+
is_primary, StringToHex(pk), txn_result.commit_ts_expired().ShortDebugString());
210+
if (is_primary) {
211+
return Status::TxnCommitTsExpired("txn commit ts expired");
212+
}
213+
} else {
214+
DINGO_LOG(FATAL) << fmt::format("[sdk.txn.{}] commit unknown txn result, is_primary({}) pk({}) response({}).",
215+
txn_id, is_primary, StringToHex(pk), response->ShortDebugString());
216+
}
217+
218+
return Status::OK();
219+
}
220+
194221
} // namespace sdk
195222

196223
} // namespace dingodb

src/sdk/transaction/txn_task/txn_commit_task.h

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -49,14 +49,15 @@ class TxnCommitTask : public TxnTask {
4949

5050
void TxnCommitRpcCallback(const Status& status, TxnCommitRpc* rpc);
5151

52+
Status ProcessTxnCommitResponse(const TxnCommitResponse* response, bool is_primary);
53+
5254
const std::vector<std::string> keys_;
5355
bool is_primary_{false};
5456
std::shared_ptr<TxnImpl> txn_impl_;
5557

5658
std::vector<StoreRpcController> controllers_;
5759
std::vector<std::unique_ptr<TxnCommitRpc>> rpcs_;
5860
std::set<std::string_view> next_keys_;
59-
bool need_retry_{false};
6061

6162
std::atomic<int> sub_tasks_count_{0};
6263
RWLock rw_lock_;

test/unit_test/sdk/transaction/test_txn_impl.cc

Lines changed: 69 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1105,6 +1105,75 @@ TEST_F(SDKTxnImplTest, RollbackSecondKeysFail) {
11051105
EXPECT_EQ(txn->TEST_IsFinishedState(), true);
11061106
}
11071107

1108+
TEST_F(SDKTxnImplTest, CommitTsExpired) {
1109+
auto txn = NewTransactionImpl(options);
1110+
1111+
EXPECT_EQ(txn->TEST_IsActiveState(), true);
1112+
1113+
{
1114+
txn->Put("a", "a");
1115+
txn->Delete("a");
1116+
1117+
txn->PutIfAbsent("d", "d");
1118+
}
1119+
int64_t original_commit_ts;
1120+
1121+
EXPECT_CALL(*rpc_client, SendRpc)
1122+
.WillOnce([&](Rpc& rpc, std::function<void()> cb) {
1123+
TxnPrewriteRpc* txn_rpc = dynamic_cast<TxnPrewriteRpc*>(&rpc);
1124+
// precommit
1125+
1126+
cb();
1127+
})
1128+
.WillOnce([&](Rpc& rpc, std::function<void()> cb) {
1129+
TxnPrewriteRpc* txn_rpc = dynamic_cast<TxnPrewriteRpc*>(&rpc);
1130+
// precommit
1131+
1132+
cb();
1133+
})
1134+
.WillOnce([&](Rpc& rpc, std::function<void()> cb) {
1135+
TxnCommitRpc* txn_rpc = dynamic_cast<TxnCommitRpc*>(&rpc);
1136+
// commit primary key
1137+
CHECK_NOTNULL(txn_rpc);
1138+
const auto* request = txn_rpc->Request();
1139+
EXPECT_TRUE(request->has_context());
1140+
EXPECT_EQ(request->start_ts(), txn->TEST_GetStartTs());
1141+
EXPECT_EQ(request->commit_ts(), txn->TEST_GetCommitTs());
1142+
original_commit_ts = request->commit_ts();
1143+
1144+
auto* response = txn_rpc->MutableResponse();
1145+
auto* commit_ts_expired = response->mutable_txn_result()->mutable_commit_ts_expired();
1146+
commit_ts_expired->mutable_key()->assign(txn->TEST_GetPrimaryKey());
1147+
1148+
cb();
1149+
})
1150+
.WillRepeatedly([&](Rpc& rpc, std::function<void()> cb) {
1151+
TxnCommitRpc* txn_rpc = dynamic_cast<TxnCommitRpc*>(&rpc);
1152+
// commit ordinary key
1153+
CHECK_NOTNULL(txn_rpc);
1154+
const auto* request = txn_rpc->Request();
1155+
EXPECT_TRUE(request->has_context());
1156+
EXPECT_EQ(request->start_ts(), txn->TEST_GetStartTs());
1157+
EXPECT_EQ(request->commit_ts(), txn->TEST_GetCommitTs());
1158+
EXPECT_NE(original_commit_ts, request->commit_ts());
1159+
1160+
cb();
1161+
});
1162+
1163+
Status s = txn->PreCommit();
1164+
EXPECT_TRUE(s.ok());
1165+
EXPECT_EQ(txn->TEST_IsPreCommittedState(), true);
1166+
1167+
s = txn->Commit();
1168+
EXPECT_TRUE(s.ok());
1169+
while (true) {
1170+
if (txn->TEST_IsFinishedState()) {
1171+
break;
1172+
}
1173+
usleep(1000);
1174+
}
1175+
}
1176+
11081177
TEST_F(SDKTxnImplTest, LockHeartbeat) {
11091178
auto txn = NewTransactionImpl(options);
11101179

0 commit comments

Comments
 (0)