mndTrans.c 47.6 KB
Newer Older
S
Shengliang Guan 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16
/*
 * Copyright (c) 2019 TAOS Data, Inc. <jhtao@taosdata.com>
 *
 * This program is free software: you can use, redistribute, and/or modify
 * it under the terms of the GNU Affero General Public License, version 3
 * or later ("AGPL"), as published by the Free Software Foundation.
 *
 * This program is distributed in the hope that it will be useful, but WITHOUT
 * ANY WARRANTY; without even the implied warranty of MERCHANTABILITY or
 * FITNESS FOR A PARTICULAR PURPOSE.
 *
 * You should have received a copy of the GNU Affero General Public License
 * along with this program. If not, see <http://www.gnu.org/licenses/>.
 */

#define _DEFAULT_SOURCE
S
Shengliang Guan 已提交
17
#include "mndTrans.h"
S
Shengliang Guan 已提交
18
#include "mndAuth.h"
L
Liu Jicong 已提交
19
#include "mndConsumer.h"
S
Shengliang Guan 已提交
20
#include "mndDb.h"
S
Shengliang Guan 已提交
21
#include "mndShow.h"
S
Shengliang Guan 已提交
22
#include "mndSync.h"
S
Shengliang Guan 已提交
23
#include "mndUser.h"
S
Shengliang Guan 已提交
24

S
Shengliang Guan 已提交
25 26 27
#define TRANS_VER_NUMBER   1
#define TRANS_ARRAY_SIZE   8
#define TRANS_RESERVE_SIZE 64
S
Shengliang Guan 已提交
28

S
Shengliang Guan 已提交
29 30 31
static SSdbRaw *mndTransActionEncode(STrans *pTrans);
static SSdbRow *mndTransActionDecode(SSdbRaw *pRaw);
static int32_t  mndTransActionInsert(SSdb *pSdb, STrans *pTrans);
S
Shengliang Guan 已提交
32
static int32_t  mndTransActionUpdate(SSdb *pSdb, STrans *OldTrans, STrans *pOld);
33
static int32_t  mndTransActionDelete(SSdb *pSdb, STrans *pTrans, bool callFunc);
S
Shengliang Guan 已提交
34

S
Shengliang Guan 已提交
35
static int32_t mndTransAppendLog(SArray *pArray, SSdbRaw *pRaw);
S
Shengliang Guan 已提交
36
static int32_t mndTransAppendAction(SArray *pArray, STransAction *pAction);
S
Shengliang Guan 已提交
37 38
static void    mndTransDropLogs(SArray *pArray);
static void    mndTransDropActions(SArray *pArray);
39
static void    mndTransDropData(STrans *pTrans);
S
Shengliang Guan 已提交
40
static int32_t mndTransExecuteLogs(SMnode *pMnode, SArray *pArray);
41
static int32_t mndTransExecuteActions(SMnode *pMnode, STrans *pTrans, SArray *pArray);
S
Shengliang Guan 已提交
42 43 44 45
static int32_t mndTransExecuteRedoLogs(SMnode *pMnode, STrans *pTrans);
static int32_t mndTransExecuteUndoLogs(SMnode *pMnode, STrans *pTrans);
static int32_t mndTransExecuteRedoActions(SMnode *pMnode, STrans *pTrans);
static int32_t mndTransExecuteUndoActions(SMnode *pMnode, STrans *pTrans);
S
Shengliang Guan 已提交
46 47 48 49 50 51 52 53 54 55 56
static int32_t mndTransExecuteCommitLogs(SMnode *pMnode, STrans *pTrans);
static bool    mndTransPerformPrepareStage(SMnode *pMnode, STrans *pTrans);
static bool    mndTransPerformRedoLogStage(SMnode *pMnode, STrans *pTrans);
static bool    mndTransPerformRedoActionStage(SMnode *pMnode, STrans *pTrans);
static bool    mndTransPerformUndoLogStage(SMnode *pMnode, STrans *pTrans);
static bool    mndTransPerformUndoActionStage(SMnode *pMnode, STrans *pTrans);
static bool    mndTransPerformCommitLogStage(SMnode *pMnode, STrans *pTrans);
static bool    mndTransPerformCommitStage(SMnode *pMnode, STrans *pTrans);
static bool    mndTransPerformRollbackStage(SMnode *pMnode, STrans *pTrans);
static bool    mndTransPerfromFinishedStage(SMnode *pMnode, STrans *pTrans);

S
Shengliang Guan 已提交
57
static void    mndTransExecute(SMnode *pMnode, STrans *pTrans);
58
static void    mndTransSendRpcRsp(SMnode *pMnode, STrans *pTrans);
S
Shengliang Guan 已提交
59 60
static int32_t mndProcessTransReq(SRpcMsg *pReq);
static int32_t mndProcessKillTransReq(SRpcMsg *pReq);
S
Shengliang Guan 已提交
61

S
Shengliang Guan 已提交
62
static int32_t mndRetrieveTrans(SRpcMsg *pReq, SShowObj *pShow, SSDataBlock *pBlock, int32_t rows);
S
Shengliang Guan 已提交
63 64
static void    mndCancelGetNextTrans(SMnode *pMnode, void *pIter);

S
Shengliang Guan 已提交
65
int32_t mndInitTrans(SMnode *pMnode) {
S
Shengliang Guan 已提交
66 67 68 69 70 71 72 73 74
  SSdbTable table = {
      .sdbType = SDB_TRANS,
      .keyType = SDB_KEY_INT32,
      .encodeFp = (SdbEncodeFp)mndTransActionEncode,
      .decodeFp = (SdbDecodeFp)mndTransActionDecode,
      .insertFp = (SdbInsertFp)mndTransActionInsert,
      .updateFp = (SdbUpdateFp)mndTransActionUpdate,
      .deleteFp = (SdbDeleteFp)mndTransActionDelete,
  };
S
Shengliang Guan 已提交
75

S
Shengliang Guan 已提交
76
  mndSetMsgHandle(pMnode, TDMT_MND_TRANS_TIMER, mndProcessTransReq);
S
Shengliang Guan 已提交
77
  mndSetMsgHandle(pMnode, TDMT_MND_KILL_TRANS, mndProcessKillTransReq);
S
Shengliang Guan 已提交
78

S
Shengliang Guan 已提交
79
  mndAddShowRetrieveHandle(pMnode, TSDB_MGMT_TABLE_TRANS, mndRetrieveTrans);
S
Shengliang Guan 已提交
80
  mndAddShowFreeIterHandle(pMnode, TSDB_MGMT_TABLE_TRANS, mndCancelGetNextTrans);
S
Shengliang Guan 已提交
81 82 83 84 85 86
  return sdbSetTable(pMnode->pSdb, table);
}

void mndCleanupTrans(SMnode *pMnode) {}

static SSdbRaw *mndTransActionEncode(STrans *pTrans) {
87 88
  terrno = TSDB_CODE_OUT_OF_MEMORY;

S
Shengliang Guan 已提交
89
  int32_t rawDataLen = sizeof(STrans) + TRANS_RESERVE_SIZE;
S
Shengliang Guan 已提交
90 91 92 93 94 95
  int32_t redoLogNum = taosArrayGetSize(pTrans->redoLogs);
  int32_t undoLogNum = taosArrayGetSize(pTrans->undoLogs);
  int32_t commitLogNum = taosArrayGetSize(pTrans->commitLogs);
  int32_t redoActionNum = taosArrayGetSize(pTrans->redoActions);
  int32_t undoActionNum = taosArrayGetSize(pTrans->undoActions);

96
  for (int32_t i = 0; i < redoLogNum; ++i) {
S
Shengliang Guan 已提交
97
    SSdbRaw *pTmp = taosArrayGetP(pTrans->redoLogs, i);
S
Shengliang Guan 已提交
98
    rawDataLen += (sdbGetRawTotalSize(pTmp) + sizeof(int32_t));
S
Shengliang Guan 已提交
99 100
  }

101
  for (int32_t i = 0; i < undoLogNum; ++i) {
S
Shengliang Guan 已提交
102
    SSdbRaw *pTmp = taosArrayGetP(pTrans->undoLogs, i);
S
Shengliang Guan 已提交
103
    rawDataLen += (sdbGetRawTotalSize(pTmp) + sizeof(int32_t));
S
Shengliang Guan 已提交
104 105
  }

106
  for (int32_t i = 0; i < commitLogNum; ++i) {
S
Shengliang Guan 已提交
107
    SSdbRaw *pTmp = taosArrayGetP(pTrans->commitLogs, i);
S
Shengliang Guan 已提交
108
    rawDataLen += (sdbGetRawTotalSize(pTmp) + sizeof(int32_t));
S
Shengliang Guan 已提交
109 110
  }

S
Shengliang Guan 已提交
111 112
  for (int32_t i = 0; i < redoActionNum; ++i) {
    STransAction *pAction = taosArrayGet(pTrans->redoActions, i);
S
Shengliang Guan 已提交
113
    rawDataLen += (sizeof(STransAction) + pAction->contLen);
S
Shengliang Guan 已提交
114 115 116 117
  }

  for (int32_t i = 0; i < undoActionNum; ++i) {
    STransAction *pAction = taosArrayGet(pTrans->undoActions, i);
S
Shengliang Guan 已提交
118
    rawDataLen += (sizeof(STransAction) + pAction->contLen);
S
Shengliang Guan 已提交
119 120
  }

S
Shengliang Guan 已提交
121
  SSdbRaw *pRaw = sdbAllocRaw(SDB_TRANS, TRANS_VER_NUMBER, rawDataLen);
S
Shengliang Guan 已提交
122
  if (pRaw == NULL) {
S
Shengliang Guan 已提交
123
    mError("trans:%d, failed to alloc raw since %s", pTrans->id, terrstr());
S
Shengliang Guan 已提交
124 125 126 127
    return NULL;
  }

  int32_t dataPos = 0;
128 129 130 131 132 133 134 135 136 137 138 139 140 141 142
  SDB_SET_INT32(pRaw, dataPos, pTrans->id, _OVER)

  ETrnStage stage = pTrans->stage;
  if (stage == TRN_STAGE_REDO_LOG || stage == TRN_STAGE_REDO_ACTION) {
    stage = TRN_STAGE_PREPARE;
  } else if (stage == TRN_STAGE_UNDO_ACTION || stage == TRN_STAGE_UNDO_LOG) {
    stage = TRN_STAGE_ROLLBACK;
  } else if (stage == TRN_STAGE_COMMIT_LOG || stage == TRN_STAGE_FINISHED) {
    stage = TRN_STAGE_COMMIT;
  } else {
  }

  SDB_SET_INT16(pRaw, dataPos, stage, _OVER)
  SDB_SET_INT16(pRaw, dataPos, pTrans->policy, _OVER)
  SDB_SET_INT16(pRaw, dataPos, pTrans->type, _OVER)
143
  SDB_SET_INT16(pRaw, dataPos, pTrans->parallel, _OVER)
144 145 146 147 148 149 150 151
  SDB_SET_INT64(pRaw, dataPos, pTrans->createdTime, _OVER)
  SDB_SET_INT64(pRaw, dataPos, pTrans->dbUid, _OVER)
  SDB_SET_BINARY(pRaw, dataPos, pTrans->dbname, TSDB_DB_FNAME_LEN, _OVER)
  SDB_SET_INT32(pRaw, dataPos, redoLogNum, _OVER)
  SDB_SET_INT32(pRaw, dataPos, undoLogNum, _OVER)
  SDB_SET_INT32(pRaw, dataPos, commitLogNum, _OVER)
  SDB_SET_INT32(pRaw, dataPos, redoActionNum, _OVER)
  SDB_SET_INT32(pRaw, dataPos, undoActionNum, _OVER)
S
Shengliang Guan 已提交
152

153
  for (int32_t i = 0; i < redoLogNum; ++i) {
S
Shengliang Guan 已提交
154
    SSdbRaw *pTmp = taosArrayGetP(pTrans->redoLogs, i);
S
Shengliang Guan 已提交
155
    int32_t  len = sdbGetRawTotalSize(pTmp);
156 157
    SDB_SET_INT32(pRaw, dataPos, len, _OVER)
    SDB_SET_BINARY(pRaw, dataPos, (void *)pTmp, len, _OVER)
S
Shengliang Guan 已提交
158 159
  }

160
  for (int32_t i = 0; i < undoLogNum; ++i) {
S
Shengliang Guan 已提交
161
    SSdbRaw *pTmp = taosArrayGetP(pTrans->undoLogs, i);
S
Shengliang Guan 已提交
162
    int32_t  len = sdbGetRawTotalSize(pTmp);
163 164
    SDB_SET_INT32(pRaw, dataPos, len, _OVER)
    SDB_SET_BINARY(pRaw, dataPos, (void *)pTmp, len, _OVER)
S
Shengliang Guan 已提交
165 166
  }

167
  for (int32_t i = 0; i < commitLogNum; ++i) {
S
Shengliang Guan 已提交
168
    SSdbRaw *pTmp = taosArrayGetP(pTrans->commitLogs, i);
S
Shengliang Guan 已提交
169
    int32_t  len = sdbGetRawTotalSize(pTmp);
170 171
    SDB_SET_INT32(pRaw, dataPos, len, _OVER)
    SDB_SET_BINARY(pRaw, dataPos, (void *)pTmp, len, _OVER)
S
Shengliang Guan 已提交
172 173
  }

S
Shengliang Guan 已提交
174 175
  for (int32_t i = 0; i < redoActionNum; ++i) {
    STransAction *pAction = taosArrayGet(pTrans->redoActions, i);
176 177 178 179 180
    SDB_SET_BINARY(pRaw, dataPos, (void *)&pAction->epSet, sizeof(SEpSet), _OVER)
    SDB_SET_INT16(pRaw, dataPos, pAction->msgType, _OVER)
    SDB_SET_INT32(pRaw, dataPos, pAction->acceptableCode, _OVER)
    SDB_SET_INT32(pRaw, dataPos, pAction->contLen, _OVER)
    SDB_SET_BINARY(pRaw, dataPos, pAction->pCont, pAction->contLen, _OVER)
S
Shengliang Guan 已提交
181 182 183 184
  }

  for (int32_t i = 0; i < undoActionNum; ++i) {
    STransAction *pAction = taosArrayGet(pTrans->undoActions, i);
185 186 187 188 189
    SDB_SET_BINARY(pRaw, dataPos, (void *)&pAction->epSet, sizeof(SEpSet), _OVER)
    SDB_SET_INT16(pRaw, dataPos, pAction->msgType, _OVER)
    SDB_SET_INT32(pRaw, dataPos, pAction->acceptableCode, _OVER)
    SDB_SET_INT32(pRaw, dataPos, pAction->contLen, _OVER)
    SDB_SET_BINARY(pRaw, dataPos, (void *)pAction->pCont, pAction->contLen, _OVER)
190 191
  }

192 193 194
  SDB_SET_INT32(pRaw, dataPos, pTrans->startFunc, _OVER)
  SDB_SET_INT32(pRaw, dataPos, pTrans->stopFunc, _OVER)
  SDB_SET_INT32(pRaw, dataPos, pTrans->paramLen, _OVER)
195
  if (pTrans->param != NULL) {
196
    SDB_SET_BINARY(pRaw, dataPos, pTrans->param, pTrans->paramLen, _OVER)
197 198
  }

199 200
  SDB_SET_RESERVE(pRaw, dataPos, TRANS_RESERVE_SIZE, _OVER)
  SDB_SET_DATALEN(pRaw, dataPos, _OVER)
201 202 203

  terrno = 0;

204
_OVER:
205 206 207 208
  if (terrno != 0) {
    mError("trans:%d, failed to encode to raw:%p len:%d since %s", pTrans->id, pRaw, dataPos, terrstr());
    sdbFreeRaw(pRaw);
    return NULL;
S
Shengliang Guan 已提交
209 210
  }

S
Shengliang Guan 已提交
211
  mTrace("trans:%d, encode to raw:%p, row:%p len:%d", pTrans->id, pRaw, pTrans, dataPos);
S
Shengliang Guan 已提交
212 213 214
  return pRaw;
}

S
Shengliang Guan 已提交
215
static SSdbRow *mndTransActionDecode(SSdbRaw *pRaw) {
216 217
  terrno = TSDB_CODE_OUT_OF_MEMORY;

218 219 220
  SSdbRow     *pRow = NULL;
  STrans      *pTrans = NULL;
  char        *pData = NULL;
221 222 223 224 225 226 227 228 229 230
  int32_t      dataLen = 0;
  int8_t       sver = 0;
  int32_t      redoLogNum = 0;
  int32_t      undoLogNum = 0;
  int32_t      commitLogNum = 0;
  int32_t      redoActionNum = 0;
  int32_t      undoActionNum = 0;
  int32_t      dataPos = 0;
  STransAction action = {0};

S
Shengliang Guan 已提交
231
  if (sdbGetRawSoftVer(pRaw, &sver) != 0) goto _OVER;
S
Shengliang Guan 已提交
232

S
Shengliang Guan 已提交
233
  if (sver != TRANS_VER_NUMBER) {
S
Shengliang Guan 已提交
234
    terrno = TSDB_CODE_SDB_INVALID_DATA_VER;
S
Shengliang Guan 已提交
235
    goto _OVER;
S
Shengliang Guan 已提交
236 237
  }

238
  pRow = sdbAllocRow(sizeof(STrans));
S
Shengliang Guan 已提交
239
  if (pRow == NULL) goto _OVER;
240 241

  pTrans = sdbGetRowObj(pRow);
S
Shengliang Guan 已提交
242
  if (pTrans == NULL) goto _OVER;
S
Shengliang Guan 已提交
243

S
Shengliang Guan 已提交
244
  SDB_GET_INT32(pRaw, dataPos, &pTrans->id, _OVER)
S
Shengliang Guan 已提交
245 246

  int16_t stage = 0;
247 248
  int16_t policy = 0;
  int16_t type = 0;
249
  int16_t parallel = 0;
S
Shengliang Guan 已提交
250
  SDB_GET_INT16(pRaw, dataPos, &stage, _OVER)
251
  SDB_GET_INT16(pRaw, dataPos, &policy, _OVER)
S
Shengliang Guan 已提交
252
  SDB_GET_INT16(pRaw, dataPos, &type, _OVER)
253
  SDB_GET_INT16(pRaw, dataPos, &parallel, _OVER)
S
Shengliang Guan 已提交
254
  pTrans->stage = stage;
255
  pTrans->policy = policy;
256
  pTrans->type = type;
257
  pTrans->parallel = parallel;
S
Shengliang Guan 已提交
258 259 260 261 262 263 264 265
  SDB_GET_INT64(pRaw, dataPos, &pTrans->createdTime, _OVER)
  SDB_GET_INT64(pRaw, dataPos, &pTrans->dbUid, _OVER)
  SDB_GET_BINARY(pRaw, dataPos, pTrans->dbname, TSDB_DB_FNAME_LEN, _OVER)
  SDB_GET_INT32(pRaw, dataPos, &redoLogNum, _OVER)
  SDB_GET_INT32(pRaw, dataPos, &undoLogNum, _OVER)
  SDB_GET_INT32(pRaw, dataPos, &commitLogNum, _OVER)
  SDB_GET_INT32(pRaw, dataPos, &redoActionNum, _OVER)
  SDB_GET_INT32(pRaw, dataPos, &undoActionNum, _OVER)
S
Shengliang Guan 已提交
266

S
Shengliang Guan 已提交
267 268 269 270 271 272
  pTrans->redoLogs = taosArrayInit(redoLogNum, sizeof(void *));
  pTrans->undoLogs = taosArrayInit(undoLogNum, sizeof(void *));
  pTrans->commitLogs = taosArrayInit(commitLogNum, sizeof(void *));
  pTrans->redoActions = taosArrayInit(redoActionNum, sizeof(STransAction));
  pTrans->undoActions = taosArrayInit(undoActionNum, sizeof(STransAction));

S
Shengliang Guan 已提交
273 274 275 276 277
  if (pTrans->redoLogs == NULL) goto _OVER;
  if (pTrans->undoLogs == NULL) goto _OVER;
  if (pTrans->commitLogs == NULL) goto _OVER;
  if (pTrans->redoActions == NULL) goto _OVER;
  if (pTrans->undoActions == NULL) goto _OVER;
S
Shengliang Guan 已提交
278

279
  for (int32_t i = 0; i < redoLogNum; ++i) {
S
Shengliang Guan 已提交
280
    SDB_GET_INT32(pRaw, dataPos, &dataLen, _OVER)
wafwerar's avatar
wafwerar 已提交
281
    pData = taosMemoryMalloc(dataLen);
S
Shengliang Guan 已提交
282
    if (pData == NULL) goto _OVER;
S
Shengliang Guan 已提交
283
    mTrace("raw:%p, is created", pData);
S
Shengliang Guan 已提交
284 285
    SDB_GET_BINARY(pRaw, dataPos, pData, dataLen, _OVER);
    if (taosArrayPush(pTrans->redoLogs, &pData) == NULL) goto _OVER;
286
    pData = NULL;
S
Shengliang Guan 已提交
287 288
  }

S
Shengliang Guan 已提交
289
  for (int32_t i = 0; i < undoLogNum; ++i) {
S
Shengliang Guan 已提交
290
    SDB_GET_INT32(pRaw, dataPos, &dataLen, _OVER)
wafwerar's avatar
wafwerar 已提交
291
    pData = taosMemoryMalloc(dataLen);
S
Shengliang Guan 已提交
292
    if (pData == NULL) goto _OVER;
S
Shengliang Guan 已提交
293
    mTrace("raw:%p, is created", pData);
S
Shengliang Guan 已提交
294 295
    SDB_GET_BINARY(pRaw, dataPos, pData, dataLen, _OVER);
    if (taosArrayPush(pTrans->undoLogs, &pData) == NULL) goto _OVER;
296
    pData = NULL;
S
Shengliang Guan 已提交
297 298 299
  }

  for (int32_t i = 0; i < commitLogNum; ++i) {
S
Shengliang Guan 已提交
300
    SDB_GET_INT32(pRaw, dataPos, &dataLen, _OVER)
wafwerar's avatar
wafwerar 已提交
301
    pData = taosMemoryMalloc(dataLen);
S
Shengliang Guan 已提交
302
    if (pData == NULL) goto _OVER;
S
Shengliang Guan 已提交
303
    mTrace("raw:%p, is created", pData);
S
Shengliang Guan 已提交
304 305
    SDB_GET_BINARY(pRaw, dataPos, pData, dataLen, _OVER);
    if (taosArrayPush(pTrans->commitLogs, &pData) == NULL) goto _OVER;
306
    pData = NULL;
S
Shengliang Guan 已提交
307 308 309
  }

  for (int32_t i = 0; i < redoActionNum; ++i) {
S
Shengliang Guan 已提交
310 311 312 313
    SDB_GET_BINARY(pRaw, dataPos, (void *)&action.epSet, sizeof(SEpSet), _OVER);
    SDB_GET_INT16(pRaw, dataPos, &action.msgType, _OVER)
    SDB_GET_INT32(pRaw, dataPos, &action.acceptableCode, _OVER)
    SDB_GET_INT32(pRaw, dataPos, &action.contLen, _OVER)
wafwerar's avatar
wafwerar 已提交
314
    action.pCont = taosMemoryMalloc(action.contLen);
S
Shengliang Guan 已提交
315 316 317
    if (action.pCont == NULL) goto _OVER;
    SDB_GET_BINARY(pRaw, dataPos, action.pCont, action.contLen, _OVER);
    if (taosArrayPush(pTrans->redoActions, &action) == NULL) goto _OVER;
318
    action.pCont = NULL;
S
Shengliang Guan 已提交
319 320 321
  }

  for (int32_t i = 0; i < undoActionNum; ++i) {
S
Shengliang Guan 已提交
322 323 324 325
    SDB_GET_BINARY(pRaw, dataPos, (void *)&action.epSet, sizeof(SEpSet), _OVER);
    SDB_GET_INT16(pRaw, dataPos, &action.msgType, _OVER)
    SDB_GET_INT32(pRaw, dataPos, &action.acceptableCode, _OVER)
    SDB_GET_INT32(pRaw, dataPos, &action.contLen, _OVER)
wafwerar's avatar
wafwerar 已提交
326
    action.pCont = taosMemoryMalloc(action.contLen);
S
Shengliang Guan 已提交
327 328 329
    if (action.pCont == NULL) goto _OVER;
    SDB_GET_BINARY(pRaw, dataPos, action.pCont, action.contLen, _OVER);
    if (taosArrayPush(pTrans->undoActions, &action) == NULL) goto _OVER;
330
    action.pCont = NULL;
S
Shengliang Guan 已提交
331 332
  }

S
Shengliang Guan 已提交
333 334 335
  SDB_GET_INT32(pRaw, dataPos, &pTrans->startFunc, _OVER)
  SDB_GET_INT32(pRaw, dataPos, &pTrans->stopFunc, _OVER)
  SDB_GET_INT32(pRaw, dataPos, &pTrans->paramLen, _OVER)
336 337
  if (pTrans->paramLen != 0) {
    pTrans->param = taosMemoryMalloc(pTrans->paramLen);
S
Shengliang Guan 已提交
338
    SDB_GET_BINARY(pRaw, dataPos, pTrans->param, pTrans->paramLen, _OVER);
339 340
  }

S
Shengliang Guan 已提交
341
  SDB_GET_RESERVE(pRaw, dataPos, TRANS_RESERVE_SIZE, _OVER)
342 343

  terrno = 0;
S
Shengliang Guan 已提交
344

S
Shengliang Guan 已提交
345
_OVER:
346 347 348
  if (terrno != 0) {
    mError("trans:%d, failed to parse from raw:%p since %s", pTrans->id, pRaw, terrstr());
    mndTransDropData(pTrans);
wafwerar's avatar
wafwerar 已提交
349 350 351
    taosMemoryFreeClear(pRow);
    taosMemoryFreeClear(pData);
    taosMemoryFreeClear(action.pCont);
S
Shengliang Guan 已提交
352 353 354
    return NULL;
  }

S
Shengliang Guan 已提交
355
  mTrace("trans:%d, decode from raw:%p, row:%p", pTrans->id, pRaw, pTrans);
S
Shengliang Guan 已提交
356 357 358
  return pRow;
}

S
Shengliang Guan 已提交
359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383
static const char *mndTransStr(ETrnStage stage) {
  switch (stage) {
    case TRN_STAGE_PREPARE:
      return "prepare";
    case TRN_STAGE_REDO_LOG:
      return "redoLog";
    case TRN_STAGE_REDO_ACTION:
      return "redoAction";
    case TRN_STAGE_COMMIT:
      return "commit";
    case TRN_STAGE_COMMIT_LOG:
      return "commitLog";
    case TRN_STAGE_UNDO_ACTION:
      return "undoAction";
    case TRN_STAGE_UNDO_LOG:
      return "undoLog";
    case TRN_STAGE_ROLLBACK:
      return "rollback";
    case TRN_STAGE_FINISHED:
      return "finished";
    default:
      return "invalid";
  }
}

S
Shengliang Guan 已提交
384 385
static const char *mndTransType(ETrnType type) {
  switch (type) {
S
Shengliang Guan 已提交
386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419
    case TRN_TYPE_CREATE_USER:
      return "create-user";
    case TRN_TYPE_ALTER_USER:
      return "alter-user";
    case TRN_TYPE_DROP_USER:
      return "drop-user";
    case TRN_TYPE_CREATE_FUNC:
      return "create-func";
    case TRN_TYPE_DROP_FUNC:
      return "drop-func";
    case TRN_TYPE_CREATE_SNODE:
      return "create-snode";
    case TRN_TYPE_DROP_SNODE:
      return "drop-snode";
    case TRN_TYPE_CREATE_QNODE:
      return "create-qnode";
    case TRN_TYPE_DROP_QNODE:
      return "drop-qnode";
    case TRN_TYPE_CREATE_BNODE:
      return "create-bnode";
    case TRN_TYPE_DROP_BNODE:
      return "drop-bnode";
    case TRN_TYPE_CREATE_MNODE:
      return "create-mnode";
    case TRN_TYPE_DROP_MNODE:
      return "drop-mnode";
    case TRN_TYPE_CREATE_TOPIC:
      return "create-topic";
    case TRN_TYPE_DROP_TOPIC:
      return "drop-topic";
    case TRN_TYPE_SUBSCRIBE:
      return "subscribe";
    case TRN_TYPE_REBALANCE:
      return "rebalance";
S
Shengliang Guan 已提交
420 421 422 423 424 425 426 427 428 429
    case TRN_TYPE_COMMIT_OFFSET:
      return "commit-offset";
    case TRN_TYPE_CREATE_STREAM:
      return "create-stream";
    case TRN_TYPE_DROP_STREAM:
      return "drop-stream";
    case TRN_TYPE_CONSUMER_LOST:
      return "consumer-lost";
    case TRN_TYPE_CONSUMER_RECOVER:
      return "consumer-recover";
S
Shengliang Guan 已提交
430 431 432 433
    case TRN_TYPE_CREATE_DNODE:
      return "create-qnode";
    case TRN_TYPE_DROP_DNODE:
      return "drop-qnode";
S
Shengliang Guan 已提交
434 435
    case TRN_TYPE_CREATE_DB:
      return "create-db";
S
Shengliang Guan 已提交
436 437 438 439 440 441 442 443 444 445 446 447 448 449
    case TRN_TYPE_ALTER_DB:
      return "alter-db";
    case TRN_TYPE_DROP_DB:
      return "drop-db";
    case TRN_TYPE_SPLIT_VGROUP:
      return "split-vgroup";
    case TRN_TYPE_MERGE_VGROUP:
      return "merge-vgroup";
    case TRN_TYPE_CREATE_STB:
      return "create-stb";
    case TRN_TYPE_ALTER_STB:
      return "alter-stb";
    case TRN_TYPE_DROP_STB:
      return "drop-stb";
S
Shengliang Guan 已提交
450 451 452 453
    case TRN_TYPE_CREATE_SMA:
      return "create-sma";
    case TRN_TYPE_DROP_SMA:
      return "drop-sma";
S
Shengliang Guan 已提交
454 455 456 457 458
    default:
      return "invalid";
  }
}

459 460 461 462 463 464 465 466 467 468 469 470 471 472
static void mndTransTestStartFunc(SMnode *pMnode, void *param, int32_t paramLen) {
  mInfo("test trans start, param:%s, len:%d", (char *)param, paramLen);
}

static void mndTransTestStopFunc(SMnode *pMnode, void *param, int32_t paramLen) {
  mInfo("test trans stop, param:%s, len:%d", (char *)param, paramLen);
}

static TransCbFp mndTransGetCbFp(ETrnFuncType ftype) {
  switch (ftype) {
    case TEST_TRANS_START_FUNC:
      return mndTransTestStartFunc;
    case TEST_TRANS_STOP_FUNC:
      return mndTransTestStopFunc;
L
Liu Jicong 已提交
473 474 475 476
    case MQ_REB_TRANS_START_FUNC:
      return mndRebCntInc;
    case MQ_REB_TRANS_STOP_FUNC:
      return mndRebCntDec;
477 478 479 480 481
    default:
      return NULL;
  }
}

S
Shengliang Guan 已提交
482
static int32_t mndTransActionInsert(SSdb *pSdb, STrans *pTrans) {
S
Shengliang Guan 已提交
483
  mTrace("trans:%d, perform insert action, row:%p stage:%s", pTrans->id, pTrans, mndTransStr(pTrans->stage));
484 485 486 487 488 489 490 491

  if (pTrans->startFunc > 0) {
    TransCbFp fp = mndTransGetCbFp(pTrans->startFunc);
    if (fp) {
      (*fp)(pSdb->pMnode, pTrans->param, pTrans->paramLen);
    }
  }

S
Shengliang Guan 已提交
492 493 494
  return 0;
}

495
static void mndTransDropData(STrans *pTrans) {
S
Shengliang Guan 已提交
496 497 498 499 500
  mndTransDropLogs(pTrans->redoLogs);
  mndTransDropLogs(pTrans->undoLogs);
  mndTransDropLogs(pTrans->commitLogs);
  mndTransDropActions(pTrans->redoActions);
  mndTransDropActions(pTrans->undoActions);
S
Shengliang Guan 已提交
501
  if (pTrans->rpcRsp != NULL) {
wafwerar's avatar
wafwerar 已提交
502
    taosMemoryFree(pTrans->rpcRsp);
S
Shengliang Guan 已提交
503 504 505
    pTrans->rpcRsp = NULL;
    pTrans->rpcRspLen = 0;
  }
506 507 508 509 510
  if (pTrans->param != NULL) {
    taosMemoryFree(pTrans->param);
    pTrans->param = NULL;
    pTrans->paramLen = 0;
  }
511
}
S
Shengliang Guan 已提交
512

513 514 515 516 517 518 519 520 521 522
static int32_t mndTransActionDelete(SSdb *pSdb, STrans *pTrans, bool callFunc) {
  mDebug("trans:%d, perform delete action, row:%p stage:%s callfunc:%d", pTrans->id, pTrans, mndTransStr(pTrans->stage),
         callFunc);
  if (pTrans->stopFunc > 0 && callFunc) {
    TransCbFp fp = mndTransGetCbFp(pTrans->stopFunc);
    if (fp) {
      (*fp)(pSdb->pMnode, pTrans->param, pTrans->paramLen);
    }
  }

523
  mndTransDropData(pTrans);
S
Shengliang Guan 已提交
524 525 526
  return 0;
}

S
Shengliang Guan 已提交
527
static int32_t mndTransActionUpdate(SSdb *pSdb, STrans *pOld, STrans *pNew) {
S
Shengliang Guan 已提交
528 529 530 531
  if (pNew->stage == TRN_STAGE_COMMIT) {
    pNew->stage = TRN_STAGE_COMMIT_LOG;
    mTrace("trans:%d, stage from %s to %s", pNew->id, mndTransStr(TRN_STAGE_COMMIT), mndTransStr(TRN_STAGE_COMMIT_LOG));
  }
532

533 534 535 536 537
  if (pNew->stage == TRN_STAGE_ROLLBACK) {
    pNew->stage = TRN_STAGE_FINISHED;
    mTrace("trans:%d, stage from %s to %s", pNew->id, mndTransStr(TRN_STAGE_ROLLBACK), mndTransStr(TRN_STAGE_FINISHED));
  }

S
Shengliang Guan 已提交
538 539
  mTrace("trans:%d, perform update action, old row:%p stage:%s, new row:%p stage:%s", pOld->id, pOld,
         mndTransStr(pOld->stage), pNew, mndTransStr(pNew->stage));
S
Shengliang Guan 已提交
540
  pOld->stage = pNew->stage;
S
Shengliang Guan 已提交
541 542 543
  return 0;
}

544
STrans *mndAcquireTrans(SMnode *pMnode, int32_t transId) {
545
  STrans *pTrans = sdbAcquire(pMnode->pSdb, SDB_TRANS, &transId);
S
Shengliang Guan 已提交
546 547 548 549
  if (pTrans == NULL) {
    terrno = TSDB_CODE_MND_TRANS_NOT_EXIST;
  }
  return pTrans;
S
Shengliang Guan 已提交
550 551
}

552
void mndReleaseTrans(SMnode *pMnode, STrans *pTrans) {
S
Shengliang Guan 已提交
553 554 555 556
  SSdb *pSdb = pMnode->pSdb;
  sdbRelease(pSdb, pTrans);
}

S
Shengliang Guan 已提交
557
STrans *mndTransCreate(SMnode *pMnode, ETrnPolicy policy, ETrnType type, const SRpcMsg *pReq) {
wafwerar's avatar
wafwerar 已提交
558
  STrans *pTrans = taosMemoryCalloc(1, sizeof(STrans));
S
Shengliang Guan 已提交
559 560 561 562 563 564
  if (pTrans == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    mError("failed to create transaction since %s", terrstr());
    return NULL;
  }

S
Shengliang Guan 已提交
565
  pTrans->id = sdbGetMaxId(pMnode->pSdb, SDB_TRANS);
S
Shengliang Guan 已提交
566 567
  pTrans->stage = TRN_STAGE_PREPARE;
  pTrans->policy = policy;
568
  pTrans->type = type;
S
Shengliang Guan 已提交
569
  pTrans->createdTime = taosGetTimestampMs();
570
  if (pReq != NULL) pTrans->rpcInfo = pReq->info;
S
Shengliang Guan 已提交
571 572 573 574 575
  pTrans->redoLogs = taosArrayInit(TRANS_ARRAY_SIZE, sizeof(void *));
  pTrans->undoLogs = taosArrayInit(TRANS_ARRAY_SIZE, sizeof(void *));
  pTrans->commitLogs = taosArrayInit(TRANS_ARRAY_SIZE, sizeof(void *));
  pTrans->redoActions = taosArrayInit(TRANS_ARRAY_SIZE, sizeof(STransAction));
  pTrans->undoActions = taosArrayInit(TRANS_ARRAY_SIZE, sizeof(STransAction));
S
Shengliang Guan 已提交
576 577 578 579 580 581 582 583

  if (pTrans->redoLogs == NULL || pTrans->undoLogs == NULL || pTrans->commitLogs == NULL ||
      pTrans->redoActions == NULL || pTrans->undoActions == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    mError("failed to create transaction since %s", terrstr());
    return NULL;
  }

584
  mDebug("trans:%d, local object is created, data:%p", pTrans->id, pTrans);
S
Shengliang Guan 已提交
585 586 587
  return pTrans;
}

S
Shengliang Guan 已提交
588
static void mndTransDropLogs(SArray *pArray) {
S
Shengliang Guan 已提交
589 590
  int32_t size = taosArrayGetSize(pArray);
  for (int32_t i = 0; i < size; ++i) {
S
Shengliang Guan 已提交
591
    SSdbRaw *pRaw = taosArrayGetP(pArray, i);
S
Shengliang Guan 已提交
592
    sdbFreeRaw(pRaw);
S
Shengliang Guan 已提交
593 594 595 596 597
  }

  taosArrayDestroy(pArray);
}

S
Shengliang Guan 已提交
598
static void mndTransDropActions(SArray *pArray) {
S
Shengliang Guan 已提交
599 600
  int32_t size = taosArrayGetSize(pArray);
  for (int32_t i = 0; i < size; ++i) {
S
Shengliang Guan 已提交
601
    STransAction *pAction = taosArrayGet(pArray, i);
wafwerar's avatar
wafwerar 已提交
602
    taosMemoryFreeClear(pAction->pCont);
S
Shengliang Guan 已提交
603 604 605 606 607
  }

  taosArrayDestroy(pArray);
}

S
Shengliang Guan 已提交
608
void mndTransDrop(STrans *pTrans) {
S
Shengliang 已提交
609 610
  if (pTrans != NULL) {
    mndTransDropData(pTrans);
611
    mDebug("trans:%d, local object is freed, data:%p", pTrans->id, pTrans);
wafwerar's avatar
wafwerar 已提交
612
    taosMemoryFreeClear(pTrans);
S
Shengliang 已提交
613
  }
S
Shengliang Guan 已提交
614 615
}

S
Shengliang Guan 已提交
616
static int32_t mndTransAppendLog(SArray *pArray, SSdbRaw *pRaw) {
S
Shengliang Guan 已提交
617
  if (pArray == NULL || pRaw == NULL) {
618
    terrno = TSDB_CODE_INVALID_PARA;
S
Shengliang Guan 已提交
619 620 621
    return -1;
  }

S
Shengliang Guan 已提交
622
  void *ptr = taosArrayPush(pArray, &pRaw);
S
Shengliang Guan 已提交
623 624 625 626 627 628 629 630
  if (ptr == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }

  return 0;
}

S
Shengliang Guan 已提交
631
int32_t mndTransAppendRedolog(STrans *pTrans, SSdbRaw *pRaw) { return mndTransAppendLog(pTrans->redoLogs, pRaw); }
S
Shengliang Guan 已提交
632

S
Shengliang Guan 已提交
633
int32_t mndTransAppendUndolog(STrans *pTrans, SSdbRaw *pRaw) { return mndTransAppendLog(pTrans->undoLogs, pRaw); }
S
Shengliang Guan 已提交
634

S
Shengliang Guan 已提交
635
int32_t mndTransAppendCommitlog(STrans *pTrans, SSdbRaw *pRaw) { return mndTransAppendLog(pTrans->commitLogs, pRaw); }
S
Shengliang Guan 已提交
636

S
Shengliang Guan 已提交
637
static int32_t mndTransAppendAction(SArray *pArray, STransAction *pAction) {
638
  void *ptr = taosArrayPush(pArray, pAction);
S
Shengliang Guan 已提交
639 640 641 642 643 644 645 646
  if (ptr == NULL) {
    terrno = TSDB_CODE_OUT_OF_MEMORY;
    return -1;
  }

  return 0;
}

S
Shengliang Guan 已提交
647
int32_t mndTransAppendRedoAction(STrans *pTrans, STransAction *pAction) {
S
Shengliang Guan 已提交
648
  return mndTransAppendAction(pTrans->redoActions, pAction);
S
Shengliang Guan 已提交
649 650
}

S
Shengliang Guan 已提交
651
int32_t mndTransAppendUndoAction(STrans *pTrans, STransAction *pAction) {
S
Shengliang Guan 已提交
652
  return mndTransAppendAction(pTrans->undoActions, pAction);
S
Shengliang Guan 已提交
653 654
}

S
Shengliang Guan 已提交
655 656 657 658 659
void mndTransSetRpcRsp(STrans *pTrans, void *pCont, int32_t contLen) {
  pTrans->rpcRsp = pCont;
  pTrans->rpcRspLen = contLen;
}

660 661 662 663 664
void mndTransSetCb(STrans *pTrans, ETrnFuncType startFunc, ETrnFuncType stopFunc, void *param, int32_t paramLen) {
  pTrans->startFunc = startFunc;
  pTrans->stopFunc = stopFunc;
  pTrans->param = param;
  pTrans->paramLen = paramLen;
665 666
}

S
Shengliang Guan 已提交
667 668 669 670 671
void mndTransSetDbInfo(STrans *pTrans, SDbObj *pDb) {
  pTrans->dbUid = pDb->uid;
  memcpy(pTrans->dbname, pDb->name, TSDB_DB_FNAME_LEN);
}

672 673
void mndTransSetExecOneByOne(STrans *pTrans) { pTrans->parallel = TRN_EXEC_ONE_BY_ONE; }

S
Shengliang Guan 已提交
674
static int32_t mndTransSync(SMnode *pMnode, STrans *pTrans) {
S
Shengliang Guan 已提交
675
  SSdbRaw *pRaw = mndTransActionEncode(pTrans);
S
Shengliang Guan 已提交
676
  if (pRaw == NULL) {
S
Shengliang Guan 已提交
677
    mError("trans:%d, failed to encode while sync trans since %s", pTrans->id, terrstr());
S
Shengliang Guan 已提交
678 679
    return -1;
  }
S
Shengliang Guan 已提交
680
  sdbSetRawStatus(pRaw, SDB_STATUS_READY);
S
Shengliang Guan 已提交
681

S
Shengliang Guan 已提交
682
  mDebug("trans:%d, sync to other nodes", pTrans->id);
S
Shengliang Guan 已提交
683
  int32_t code = mndSyncPropose(pMnode, pRaw, pTrans->id);
S
Shengliang Guan 已提交
684 685 686
  if (code != 0) {
    mError("trans:%d, failed to sync since %s", pTrans->id, terrstr());
    sdbFreeRaw(pRaw);
S
Shengliang Guan 已提交
687 688 689
    return -1;
  }

690
  sdbFreeRaw(pRaw);
S
Shengliang Guan 已提交
691
  mDebug("trans:%d, sync finished", pTrans->id);
S
Shengliang Guan 已提交
692 693 694
  return 0;
}

S
Shengliang Guan 已提交
695
static bool mndIsBasicTrans(STrans *pTrans) {
696
  return pTrans->type > TRN_TYPE_BASIC_SCOPE && pTrans->type < TRN_TYPE_BASIC_SCOPE_END;
S
Shengliang Guan 已提交
697 698 699
}

static bool mndIsGlobalTrans(STrans *pTrans) {
700
  return pTrans->type > TRN_TYPE_GLOBAL_SCOPE && pTrans->type < TRN_TYPE_GLOBAL_SCOPE_END;
S
Shengliang Guan 已提交
701 702 703
}

static bool mndIsDbTrans(STrans *pTrans) {
704
  return pTrans->type > TRN_TYPE_DB_SCOPE && pTrans->type < TRN_TYPE_DB_SCOPE_END;
S
Shengliang Guan 已提交
705 706 707
}

static bool mndIsStbTrans(STrans *pTrans) {
708
  return pTrans->type > TRN_TYPE_STB_SCOPE && pTrans->type < TRN_TYPE_STB_SCOPE_END;
S
Shengliang Guan 已提交
709 710
}

711
static bool mndCheckTransConflict(SMnode *pMnode, STrans *pNewTrans) {
S
Shengliang Guan 已提交
712
  STrans *pTrans = NULL;
713
  void   *pIter = NULL;
714
  bool    conflict = false;
715

716
  if (mndIsBasicTrans(pNewTrans)) return conflict;
S
Shengliang Guan 已提交
717 718 719 720 721 722 723 724

  while (1) {
    pIter = sdbFetch(pMnode->pSdb, SDB_TRANS, pIter, (void **)&pTrans);
    if (pIter == NULL) break;

    if (mndIsGlobalTrans(pNewTrans)) {
      if (mndIsDbTrans(pTrans) || mndIsStbTrans(pTrans)) {
        mError("trans:%d, can't execute since trans:%d in progress db:%s", pNewTrans->id, pTrans->id, pTrans->dbname);
725
        conflict = true;
S
Shengliang Guan 已提交
726
      } else {
S
Shengliang Guan 已提交
727 728 729
      }
    }

S
Shengliang Guan 已提交
730
    else if (mndIsDbTrans(pNewTrans)) {
S
Shengliang Guan 已提交
731 732
      if (mndIsGlobalTrans(pTrans)) {
        mError("trans:%d, can't execute since trans:%d in progress", pNewTrans->id, pTrans->id);
733
        conflict = true;
S
Shengliang Guan 已提交
734
      } else if (mndIsDbTrans(pTrans) || mndIsStbTrans(pTrans)) {
S
Shengliang Guan 已提交
735 736
        if (pNewTrans->dbUid == pTrans->dbUid) {
          mError("trans:%d, can't execute since trans:%d in progress db:%s", pNewTrans->id, pTrans->id, pTrans->dbname);
737
          conflict = true;
S
Shengliang Guan 已提交
738
        }
S
Shengliang Guan 已提交
739
      } else {
S
Shengliang Guan 已提交
740 741 742
      }
    }

S
Shengliang Guan 已提交
743
    else if (mndIsStbTrans(pNewTrans)) {
S
Shengliang Guan 已提交
744 745
      if (mndIsGlobalTrans(pTrans)) {
        mError("trans:%d, can't execute since trans:%d in progress", pNewTrans->id, pTrans->id);
746
        conflict = true;
S
Shengliang Guan 已提交
747
      } else if (mndIsDbTrans(pTrans)) {
S
Shengliang Guan 已提交
748 749
        if (pNewTrans->dbUid == pTrans->dbUid) {
          mError("trans:%d, can't execute since trans:%d in progress db:%s", pNewTrans->id, pTrans->id, pTrans->dbname);
750
          conflict = true;
S
Shengliang Guan 已提交
751
        }
S
Shengliang Guan 已提交
752
      } else {
S
Shengliang Guan 已提交
753 754 755 756 757 758 759 760
      }
    }

    sdbRelease(pMnode->pSdb, pTrans);
  }

  sdbCancelFetch(pMnode->pSdb, pIter);
  sdbRelease(pMnode->pSdb, pTrans);
761
  return conflict;
S
Shengliang Guan 已提交
762 763
}

S
Shengliang Guan 已提交
764
int32_t mndTransPrepare(SMnode *pMnode, STrans *pTrans) {
765 766
  if (mndCheckTransConflict(pMnode, pTrans)) {
    terrno = TSDB_CODE_MND_TRANS_CONFLICT;
S
Shengliang Guan 已提交
767 768 769 770
    mError("trans:%d, failed to prepare since %s", pTrans->id, terrstr());
    return -1;
  }

771 772 773 774 775 776
  if (taosArrayGetSize(pTrans->commitLogs) <= 0) {
    terrno = TSDB_CODE_MND_TRANS_CLOG_IS_NULL;
    mError("trans:%d, failed to prepare since %s", pTrans->id, terrstr());
    return -1;
  }

S
Shengliang Guan 已提交
777 778 779 780 781 782 783
  mDebug("trans:%d, prepare transaction", pTrans->id);
  if (mndTransSync(pMnode, pTrans) != 0) {
    mError("trans:%d, failed to prepare since %s", pTrans->id, terrstr());
    return -1;
  }
  mDebug("trans:%d, prepare finished", pTrans->id);

S
Shengliang Guan 已提交
784 785
  STrans *pNew = mndAcquireTrans(pMnode, pTrans->id);
  if (pNew == NULL) {
S
Shengliang Guan 已提交
786
    mError("trans:%d, failed to read from sdb since %s", pTrans->id, terrstr());
S
Shengliang Guan 已提交
787 788 789
    return -1;
  }

S
Shengliang Guan 已提交
790
  pNew->rpcInfo = pTrans->rpcInfo;
S
Shengliang Guan 已提交
791 792 793 794 795
  pNew->rpcRsp = pTrans->rpcRsp;
  pNew->rpcRspLen = pTrans->rpcRspLen;
  pTrans->rpcRsp = NULL;
  pTrans->rpcRspLen = 0;

S
Shengliang Guan 已提交
796 797
  mndTransExecute(pMnode, pNew);
  mndReleaseTrans(pMnode, pNew);
S
Shengliang Guan 已提交
798 799 800
  return 0;
}

S
Shengliang Guan 已提交
801 802
static int32_t mndTransCommit(SMnode *pMnode, STrans *pTrans) {
  if (taosArrayGetSize(pTrans->commitLogs) == 0 && taosArrayGetSize(pTrans->redoActions) == 0) return 0;
S
Shengliang Guan 已提交
803

S
Shengliang Guan 已提交
804 805 806
  mDebug("trans:%d, commit transaction", pTrans->id);
  if (mndTransSync(pMnode, pTrans) != 0) {
    mError("trans:%d, failed to commit since %s", pTrans->id, terrstr());
S
Shengliang Guan 已提交
807
    return -1;
S
Shengliang Guan 已提交
808 809
  }
  mDebug("trans:%d, commit finished", pTrans->id);
S
Shengliang Guan 已提交
810 811 812
  return 0;
}

S
Shengliang Guan 已提交
813
static int32_t mndTransRollback(SMnode *pMnode, STrans *pTrans) {
S
Shengliang Guan 已提交
814
  mDebug("trans:%d, rollback transaction", pTrans->id);
S
Shengliang Guan 已提交
815 816
  if (mndTransSync(pMnode, pTrans) != 0) {
    mError("trans:%d, failed to rollback since %s", pTrans->id, terrstr());
S
Shengliang Guan 已提交
817
    return -1;
S
Shengliang Guan 已提交
818
  }
S
Shengliang Guan 已提交
819 820
  mDebug("trans:%d, rollback finished", pTrans->id);
  return 0;
S
Shengliang Guan 已提交
821
}
S
Shengliang Guan 已提交
822

823
static void mndTransSendRpcRsp(SMnode *pMnode, STrans *pTrans) {
824 825
  bool    sendRsp = false;
  int32_t code = pTrans->code;
826 827 828 829 830 831 832 833

  if (pTrans->stage == TRN_STAGE_FINISHED) {
    sendRsp = true;
  }

  if (pTrans->policy == TRN_POLICY_ROLLBACK) {
    if (pTrans->stage == TRN_STAGE_UNDO_LOG || pTrans->stage == TRN_STAGE_UNDO_ACTION ||
        pTrans->stage == TRN_STAGE_ROLLBACK) {
834
      if (code == 0) code = TSDB_CODE_MND_TRANS_UNKNOW_ERROR;
835
      sendRsp = true;
836
    }
837
  } else {
S
Shengliang Guan 已提交
838
    if (pTrans->stage == TRN_STAGE_REDO_ACTION && pTrans->failedTimes > 6) {
839
      if (code == 0) code = TSDB_CODE_MND_TRANS_UNKNOW_ERROR;
840
      sendRsp = true;
841
    }
S
Shengliang Guan 已提交
842
  }
843

S
Shengliang Guan 已提交
844
  if (sendRsp && pTrans->rpcInfo.handle != NULL) {
S
Shengliang Guan 已提交
845 846 847 848
    void *rpcCont = rpcMallocCont(pTrans->rpcRspLen);
    if (rpcCont != NULL) {
      memcpy(rpcCont, pTrans->rpcRsp, pTrans->rpcRspLen);
    }
wafwerar's avatar
wafwerar 已提交
849
    taosMemoryFree(pTrans->rpcRsp);
S
Shengliang Guan 已提交
850

851
    mDebug("trans:%d, send rsp, code:0x%x stage:%d app:%p", pTrans->id, code, pTrans->stage, pTrans->rpcInfo.ahandle);
S
Shengliang Guan 已提交
852
    SRpcMsg rspMsg = {
853
        .code = code,
S
Shengliang Guan 已提交
854 855
        .pCont = rpcCont,
        .contLen = pTrans->rpcRspLen,
856
        .info = pTrans->rpcInfo,
S
Shengliang Guan 已提交
857
    };
S
shm  
Shengliang Guan 已提交
858
    tmsgSendRsp(&rspMsg);
S
Shengliang Guan 已提交
859
    pTrans->rpcInfo.handle = NULL;
S
Shengliang Guan 已提交
860 861
    pTrans->rpcRsp = NULL;
    pTrans->rpcRspLen = 0;
862
  }
S
Shengliang Guan 已提交
863 864
}

S
Shengliang Guan 已提交
865 866 867
void mndTransProcessRsp(SRpcMsg *pRsp) {
  SMnode *pMnode = pRsp->info.node;
  int64_t signature = (int64_t)(pRsp->info.ahandle);
S
Shengliang Guan 已提交
868 869
  int32_t transId = (int32_t)(signature >> 32);
  int32_t action = (int32_t)((signature << 32) >> 32);
870 871 872 873

  STrans *pTrans = mndAcquireTrans(pMnode, transId);
  if (pTrans == NULL) {
    mError("trans:%d, failed to get transId from vnode rsp since %s", transId, terrstr());
S
Shengliang Guan 已提交
874
    goto _OVER;
875 876 877
  }

  SArray *pArray = NULL;
S
Shengliang Guan 已提交
878
  if (pTrans->stage == TRN_STAGE_REDO_ACTION) {
879
    pArray = pTrans->redoActions;
S
Shengliang Guan 已提交
880
  } else if (pTrans->stage == TRN_STAGE_UNDO_ACTION) {
881 882
    pArray = pTrans->undoActions;
  } else {
S
Shengliang Guan 已提交
883
    mError("trans:%d, invalid trans stage:%d while recv action rsp", pTrans->id, pTrans->stage);
S
Shengliang Guan 已提交
884
    goto _OVER;
885 886 887
  }

  if (pArray == NULL) {
S
Shengliang Guan 已提交
888
    mError("trans:%d, invalid trans stage:%d", transId, pTrans->stage);
S
Shengliang Guan 已提交
889
    goto _OVER;
890 891 892
  }

  int32_t actionNum = taosArrayGetSize(pTrans->redoActions);
S
Shengliang Guan 已提交
893
  if (action < 0 || action >= actionNum) {
894
    mError("trans:%d, invalid action:%d", transId, action);
S
Shengliang Guan 已提交
895
    goto _OVER;
896 897 898 899 900
  }

  STransAction *pAction = taosArrayGet(pArray, action);
  if (pAction != NULL) {
    pAction->msgReceived = 1;
S
Shengliang Guan 已提交
901
    pAction->errCode = pRsp->code;
S
Shengliang Guan 已提交
902 903 904
    if (pAction->errCode != 0) {
      tstrncpy(pTrans->lastError, tstrerror(pAction->errCode), TSDB_TRANS_ERROR_LEN);
    }
905 906
  }

S
Shengliang Guan 已提交
907
  mDebug("trans:%d, action:%d response is received, code:0x%x, accept:0x%04x", transId, action, pRsp->code,
908
         pAction->acceptableCode);
909 910
  mndTransExecute(pMnode, pTrans);

S
Shengliang Guan 已提交
911
_OVER:
912 913 914
  mndReleaseTrans(pMnode, pTrans);
}

S
Shengliang Guan 已提交
915
static int32_t mndTransExecuteLogs(SMnode *pMnode, SArray *pArray) {
916
  SSdb   *pSdb = pMnode->pSdb;
S
Shengliang Guan 已提交
917 918
  int32_t arraySize = taosArrayGetSize(pArray);

S
Shengliang Guan 已提交
919 920
  if (arraySize == 0) return 0;

921
  int32_t code = 0;
S
Shengliang Guan 已提交
922
  for (int32_t i = 0; i < arraySize; ++i) {
923
    SSdbRaw *pRaw = taosArrayGetP(pArray, i);
924 925
    if (sdbWriteWithoutFree(pSdb, pRaw) != 0) {
      code = ((terrno != 0) ? terrno : -1);
S
Shengliang Guan 已提交
926 927 928
    }
  }

929 930
  terrno = code;
  return code;
S
Shengliang Guan 已提交
931 932
}

S
Shengliang Guan 已提交
933
static int32_t mndTransExecuteRedoLogs(SMnode *pMnode, STrans *pTrans) {
934 935 936 937 938
  int32_t code = mndTransExecuteLogs(pMnode, pTrans->redoLogs);
  if (code != 0) {
    mError("failed to execute redoLogs since %s", terrstr());
  }
  return code;
S
Shengliang Guan 已提交
939 940 941
}

static int32_t mndTransExecuteUndoLogs(SMnode *pMnode, STrans *pTrans) {
942 943 944 945 946 947
  int32_t code = mndTransExecuteLogs(pMnode, pTrans->undoLogs);
  if (code != 0) {
    mError("failed to execute undoLogs since %s, return success", terrstr());
  }

  return 0;  // return success in any case
S
Shengliang Guan 已提交
948 949 950
}

static int32_t mndTransExecuteCommitLogs(SMnode *pMnode, STrans *pTrans) {
951 952 953 954 955
  int32_t code = mndTransExecuteLogs(pMnode, pTrans->commitLogs);
  if (code != 0) {
    mError("failed to execute commitLogs since %s", terrstr());
  }
  return code;
S
Shengliang Guan 已提交
956
}
S
Shengliang Guan 已提交
957

S
Shengliang Guan 已提交
958 959 960 961 962 963 964 965 966 967 968
static void mndTransResetActions(SMnode *pMnode, STrans *pTrans, SArray *pArray) {
  int32_t numOfActions = taosArrayGetSize(pArray);

  for (int32_t action = 0; action < numOfActions; ++action) {
    STransAction *pAction = taosArrayGet(pArray, action);
    if (pAction == NULL) continue;
    if (pAction->msgSent && pAction->msgReceived && pAction->errCode == 0) continue;

    pAction->msgSent = 0;
    pAction->msgReceived = 0;
    pAction->errCode = 0;
S
Shengliang Guan 已提交
969
    mDebug("trans:%d, action:%d execute status is reset", pTrans->id, action);
S
Shengliang Guan 已提交
970 971 972 973
  }
}

static int32_t mndTransSendActionMsg(SMnode *pMnode, STrans *pTrans, SArray *pArray) {
974 975 976 977 978
  int32_t numOfActions = taosArrayGetSize(pArray);

  for (int32_t action = 0; action < numOfActions; ++action) {
    STransAction *pAction = taosArrayGet(pArray, action);
    if (pAction == NULL) continue;
979 980 981 982 983 984 985 986 987 988 989 990

    if (pAction->msgSent) {
      if (pAction->msgReceived) {
        continue;
      } else {
        if (pTrans->parallel == TRN_EXEC_ONE_BY_ONE) {
          break;
        } else {
          continue;
        }
      }
    }
S
Shengliang Guan 已提交
991

992 993 994 995
    int64_t signature = pTrans->id;
    signature = (signature << 32);
    signature += action;

S
Shengliang Guan 已提交
996
    SRpcMsg rpcMsg = {.msgType = pAction->msgType, .contLen = pAction->contLen, .info.ahandle = (void *)signature};
S
Shengliang Guan 已提交
997 998 999 1000
    rpcMsg.pCont = rpcMallocCont(pAction->contLen);
    if (rpcMsg.pCont == NULL) {
      terrno = TSDB_CODE_OUT_OF_MEMORY;
      return -1;
S
Shengliang Guan 已提交
1001
    }
S
Shengliang Guan 已提交
1002
    memcpy(rpcMsg.pCont, pAction->pCont, pAction->contLen);
1003

1004
    if (tmsgSendReq(&pAction->epSet, &rpcMsg) == 0) {
1005 1006
      mDebug("trans:%d, action:%d is sent to %s:%u", pTrans->id, action, pAction->epSet.eps[pAction->epSet.inUse].fqdn,
             pAction->epSet.eps[pAction->epSet.inUse].port);
S
Shengliang Guan 已提交
1007 1008 1009
      pAction->msgSent = 1;
      pAction->msgReceived = 0;
      pAction->errCode = 0;
1010 1011 1012
      if (pTrans->parallel == TRN_EXEC_ONE_BY_ONE) {
        break;
      }
S
Shengliang Guan 已提交
1013
    } else {
1014 1015 1016
      pAction->msgSent = 0;
      pAction->msgReceived = 0;
      pAction->errCode = terrno;
S
shm  
Shengliang Guan 已提交
1017
      mError("trans:%d, action:%d not send since %s", pTrans->id, action, terrstr());
S
Shengliang Guan 已提交
1018 1019
      return -1;
    }
S
Shengliang Guan 已提交
1020 1021
  }

S
Shengliang Guan 已提交
1022 1023 1024 1025 1026 1027 1028 1029 1030 1031 1032
  return 0;
}

static int32_t mndTransExecuteActions(SMnode *pMnode, STrans *pTrans, SArray *pArray) {
  int32_t numOfActions = taosArrayGetSize(pArray);
  if (numOfActions == 0) return 0;

  if (mndTransSendActionMsg(pMnode, pTrans, pArray) != 0) {
    return -1;
  }

S
Shengliang Guan 已提交
1033 1034
  int32_t numOfReceived = 0;
  int32_t errCode = 0;
1035 1036 1037 1038
  for (int32_t action = 0; action < numOfActions; ++action) {
    STransAction *pAction = taosArrayGet(pArray, action);
    if (pAction == NULL) continue;
    if (pAction->msgSent && pAction->msgReceived) {
S
Shengliang Guan 已提交
1039
      numOfReceived++;
1040
      if (pAction->errCode != 0 && pAction->errCode != pAction->acceptableCode) {
S
Shengliang Guan 已提交
1041
        errCode = pAction->errCode;
1042 1043 1044 1045
      }
    }
  }

S
Shengliang Guan 已提交
1046
  if (numOfReceived == numOfActions) {
S
Shengliang Guan 已提交
1047 1048 1049 1050
    if (errCode == 0) {
      mDebug("trans:%d, all %d actions execute successfully", pTrans->id, numOfActions);
      return 0;
    } else {
S
Shengliang Guan 已提交
1051
      mError("trans:%d, all %d actions executed, code:0x%x", pTrans->id, numOfActions, errCode & 0XFFFF);
S
Shengliang Guan 已提交
1052 1053 1054 1055
      mndTransResetActions(pMnode, pTrans, pArray);
      terrno = errCode;
      return errCode;
    }
1056
  } else {
S
Shengliang Guan 已提交
1057
    mDebug("trans:%d, %d of %d actions executed", pTrans->id, numOfReceived, numOfActions);
S
Shengliang Guan 已提交
1058
    return TSDB_CODE_ACTION_IN_PROGRESS;
1059
  }
S
Shengliang Guan 已提交
1060 1061
}

S
Shengliang Guan 已提交
1062
static int32_t mndTransExecuteRedoActions(SMnode *pMnode, STrans *pTrans) {
1063
  int32_t code = mndTransExecuteActions(pMnode, pTrans, pTrans->redoActions);
S
Shengliang Guan 已提交
1064
  if (code != 0 && code != TSDB_CODE_ACTION_IN_PROGRESS) {
1065
    mError("failed to execute redoActions since:%s, code:0x%x", terrstr(), terrno);
1066 1067
  }
  return code;
S
Shengliang Guan 已提交
1068
}
S
Shengliang Guan 已提交
1069

S
Shengliang Guan 已提交
1070
static int32_t mndTransExecuteUndoActions(SMnode *pMnode, STrans *pTrans) {
1071
  int32_t code = mndTransExecuteActions(pMnode, pTrans, pTrans->undoActions);
S
Shengliang Guan 已提交
1072
  if (code != 0 && code != TSDB_CODE_ACTION_IN_PROGRESS) {
1073 1074 1075
    mError("failed to execute undoActions since %s", terrstr());
  }
  return code;
S
Shengliang Guan 已提交
1076
}
S
Shengliang Guan 已提交
1077

S
Shengliang Guan 已提交
1078 1079 1080 1081 1082 1083 1084 1085 1086
static bool mndTransPerformPrepareStage(SMnode *pMnode, STrans *pTrans) {
  bool continueExec = true;
  pTrans->stage = TRN_STAGE_REDO_LOG;
  mDebug("trans:%d, stage from prepare to redoLog", pTrans->id);
  return continueExec;
}

static bool mndTransPerformRedoLogStage(SMnode *pMnode, STrans *pTrans) {
  bool    continueExec = true;
S
Shengliang Guan 已提交
1087
  int32_t code = mndTransExecuteRedoLogs(pMnode, pTrans);
S
Shengliang Guan 已提交
1088

S
Shengliang Guan 已提交
1089
  if (code == 0) {
S
Shengliang Guan 已提交
1090 1091 1092
    pTrans->code = 0;
    pTrans->stage = TRN_STAGE_REDO_ACTION;
    mDebug("trans:%d, stage from redoLog to redoAction", pTrans->id);
S
Shengliang Guan 已提交
1093
  } else {
S
Shengliang Guan 已提交
1094 1095
    pTrans->code = terrno;
    pTrans->stage = TRN_STAGE_UNDO_LOG;
1096
    mError("trans:%d, stage from redoLog to undoLog since %s", pTrans->id, terrstr());
S
Shengliang Guan 已提交
1097
  }
S
Shengliang Guan 已提交
1098 1099

  return continueExec;
S
Shengliang Guan 已提交
1100 1101
}

S
Shengliang Guan 已提交
1102
static bool mndTransPerformRedoActionStage(SMnode *pMnode, STrans *pTrans) {
S
Shengliang Guan 已提交
1103
  if (!pMnode->deploy && !mndIsMaster(pMnode)) return false;
1104

S
Shengliang Guan 已提交
1105
  bool    continueExec = true;
S
Shengliang Guan 已提交
1106
  int32_t code = mndTransExecuteRedoActions(pMnode, pTrans);
S
Shengliang Guan 已提交
1107 1108

  if (code == 0) {
S
Shengliang Guan 已提交
1109
    pTrans->code = 0;
S
Shengliang Guan 已提交
1110
    pTrans->stage = TRN_STAGE_COMMIT;
S
Shengliang Guan 已提交
1111 1112
    mDebug("trans:%d, stage from redoAction to commit", pTrans->id);
    continueExec = true;
S
Shengliang Guan 已提交
1113
  } else if (code == TSDB_CODE_ACTION_IN_PROGRESS) {
S
Shengliang Guan 已提交
1114 1115
    mDebug("trans:%d, stage keep on redoAction since %s", pTrans->id, tstrerror(code));
    continueExec = false;
S
Shengliang Guan 已提交
1116
  } else {
S
Shengliang Guan 已提交
1117
    pTrans->code = terrno;
S
Shengliang Guan 已提交
1118
    if (pTrans->policy == TRN_POLICY_ROLLBACK) {
S
Shengliang Guan 已提交
1119 1120 1121
      pTrans->stage = TRN_STAGE_UNDO_ACTION;
      mError("trans:%d, stage from redoAction to undoAction since %s", pTrans->id, terrstr());
      continueExec = true;
S
Shengliang Guan 已提交
1122
    } else {
S
Shengliang Guan 已提交
1123 1124 1125
      pTrans->failedTimes++;
      mError("trans:%d, stage keep on redoAction since %s, failedTimes:%d", pTrans->id, terrstr(), pTrans->failedTimes);
      continueExec = false;
S
Shengliang Guan 已提交
1126 1127 1128
    }
  }

S
Shengliang Guan 已提交
1129
  return continueExec;
S
Shengliang Guan 已提交
1130 1131
}

S
Shengliang Guan 已提交
1132 1133 1134 1135 1136 1137 1138 1139 1140 1141 1142 1143
static bool mndTransPerformCommitStage(SMnode *pMnode, STrans *pTrans) {
  bool    continueExec = true;
  int32_t code = mndTransCommit(pMnode, pTrans);

  if (code == 0) {
    pTrans->code = 0;
    pTrans->stage = TRN_STAGE_COMMIT_LOG;
    mDebug("trans:%d, stage from commit to commitLog", pTrans->id);
    continueExec = true;
  } else {
    pTrans->code = terrno;
    if (pTrans->policy == TRN_POLICY_ROLLBACK) {
1144 1145
      pTrans->stage = TRN_STAGE_UNDO_ACTION;
      mError("trans:%d, stage from commit to undoAction since %s, failedTimes:%d", pTrans->id, terrstr(),
S
Shengliang Guan 已提交
1146 1147 1148 1149 1150 1151 1152 1153 1154 1155
             pTrans->failedTimes);
      continueExec = true;
    } else {
      pTrans->failedTimes++;
      mError("trans:%d, stage keep on commit since %s, failedTimes:%d", pTrans->id, terrstr(), pTrans->failedTimes);
      continueExec = false;
    }
  }

  return continueExec;
S
Shengliang Guan 已提交
1156 1157
}

S
Shengliang Guan 已提交
1158 1159 1160
static bool mndTransPerformCommitLogStage(SMnode *pMnode, STrans *pTrans) {
  bool    continueExec = true;
  int32_t code = mndTransExecuteCommitLogs(pMnode, pTrans);
S
Shengliang Guan 已提交
1161 1162

  if (code == 0) {
S
Shengliang Guan 已提交
1163 1164 1165 1166
    pTrans->code = 0;
    pTrans->stage = TRN_STAGE_FINISHED;
    mDebug("trans:%d, stage from commitLog to finished", pTrans->id);
    continueExec = true;
S
Shengliang Guan 已提交
1167
  } else {
S
Shengliang Guan 已提交
1168 1169
    pTrans->code = terrno;
    pTrans->failedTimes++;
1170
    mError("trans:%d, stage keep on commitLog since %s, failedTimes:%d", pTrans->id, terrstr(), pTrans->failedTimes);
S
Shengliang Guan 已提交
1171 1172 1173 1174 1175 1176 1177 1178 1179 1180 1181
    continueExec = false;
  }

  return continueExec;
}

static bool mndTransPerformUndoLogStage(SMnode *pMnode, STrans *pTrans) {
  bool    continueExec = true;
  int32_t code = mndTransExecuteUndoLogs(pMnode, pTrans);

  if (code == 0) {
S
Shengliang Guan 已提交
1182
    pTrans->stage = TRN_STAGE_ROLLBACK;
S
Shengliang Guan 已提交
1183 1184 1185
    mDebug("trans:%d, stage from undoLog to rollback", pTrans->id);
    continueExec = true;
  } else {
S
Shengliang Guan 已提交
1186
    mError("trans:%d, stage keep on undoLog since %s", pTrans->id, terrstr());
S
Shengliang Guan 已提交
1187 1188 1189 1190 1191 1192 1193
    continueExec = false;
  }

  return continueExec;
}

static bool mndTransPerformUndoActionStage(SMnode *pMnode, STrans *pTrans) {
S
Shengliang Guan 已提交
1194
  if (!pMnode->deploy && !mndIsMaster(pMnode)) return false;
1195

S
Shengliang Guan 已提交
1196 1197 1198 1199
  bool    continueExec = true;
  int32_t code = mndTransExecuteUndoActions(pMnode, pTrans);

  if (code == 0) {
1200
    pTrans->stage = TRN_STAGE_UNDO_LOG;
S
Shengliang Guan 已提交
1201 1202
    mDebug("trans:%d, stage from undoAction to undoLog", pTrans->id);
    continueExec = true;
S
Shengliang Guan 已提交
1203
  } else if (code == TSDB_CODE_ACTION_IN_PROGRESS) {
S
Shengliang Guan 已提交
1204
    mDebug("trans:%d, stage keep on undoAction since %s", pTrans->id, tstrerror(code));
S
Shengliang Guan 已提交
1205 1206 1207
    continueExec = false;
  } else {
    pTrans->failedTimes++;
1208
    mError("trans:%d, stage keep on undoAction since %s, failedTimes:%d", pTrans->id, terrstr(), pTrans->failedTimes);
S
Shengliang Guan 已提交
1209 1210 1211 1212 1213 1214 1215 1216 1217 1218 1219 1220 1221 1222 1223 1224
    continueExec = false;
  }

  return continueExec;
}

static bool mndTransPerformRollbackStage(SMnode *pMnode, STrans *pTrans) {
  bool    continueExec = true;
  int32_t code = mndTransRollback(pMnode, pTrans);

  if (code == 0) {
    pTrans->stage = TRN_STAGE_FINISHED;
    mDebug("trans:%d, stage from rollback to finished", pTrans->id);
    continueExec = true;
  } else {
    pTrans->failedTimes++;
1225
    mError("trans:%d, stage keep on rollback since %s, failedTimes:%d", pTrans->id, terrstr(), pTrans->failedTimes);
S
Shengliang Guan 已提交
1226
    continueExec = false;
S
Shengliang Guan 已提交
1227 1228
  }

S
Shengliang Guan 已提交
1229 1230 1231 1232 1233 1234 1235 1236
  return continueExec;
}

static bool mndTransPerfromFinishedStage(SMnode *pMnode, STrans *pTrans) {
  bool continueExec = false;

  SSdbRaw *pRaw = mndTransActionEncode(pTrans);
  if (pRaw == NULL) {
S
Shengliang Guan 已提交
1237
    mError("trans:%d, failed to encode while finish trans since %s", pTrans->id, terrstr());
S
Shengliang Guan 已提交
1238 1239 1240 1241 1242 1243 1244 1245
  }
  sdbSetRawStatus(pRaw, SDB_STATUS_DROPPED);

  int32_t code = sdbWrite(pMnode->pSdb, pRaw);
  if (code != 0) {
    mError("trans:%d, failed to write sdb since %s", pTrans->id, terrstr());
  }

S
Shengliang Guan 已提交
1246
  mDebug("trans:%d, finished, code:0x%x, failedTimes:%d", pTrans->id, pTrans->code, pTrans->failedTimes);
1247

S
Shengliang Guan 已提交
1248
  return continueExec;
S
Shengliang Guan 已提交
1249
}
S
Shengliang Guan 已提交
1250

S
Shengliang Guan 已提交
1251
static void mndTransExecute(SMnode *pMnode, STrans *pTrans) {
S
Shengliang Guan 已提交
1252
  bool continueExec = true;
S
Shengliang Guan 已提交
1253

S
Shengliang Guan 已提交
1254
  while (continueExec) {
S
Shengliang Guan 已提交
1255
    pTrans->lastExecTime = taosGetTimestampMs();
S
Shengliang Guan 已提交
1256 1257
    switch (pTrans->stage) {
      case TRN_STAGE_PREPARE:
S
Shengliang Guan 已提交
1258 1259 1260 1261 1262 1263 1264 1265 1266 1267
        continueExec = mndTransPerformPrepareStage(pMnode, pTrans);
        break;
      case TRN_STAGE_REDO_LOG:
        continueExec = mndTransPerformRedoLogStage(pMnode, pTrans);
        break;
      case TRN_STAGE_REDO_ACTION:
        continueExec = mndTransPerformRedoActionStage(pMnode, pTrans);
        break;
      case TRN_STAGE_UNDO_LOG:
        continueExec = mndTransPerformUndoLogStage(pMnode, pTrans);
S
Shengliang Guan 已提交
1268
        break;
S
Shengliang Guan 已提交
1269 1270 1271 1272 1273
      case TRN_STAGE_UNDO_ACTION:
        continueExec = mndTransPerformUndoActionStage(pMnode, pTrans);
        break;
      case TRN_STAGE_COMMIT_LOG:
        continueExec = mndTransPerformCommitLogStage(pMnode, pTrans);
S
Shengliang Guan 已提交
1274 1275
        break;
      case TRN_STAGE_COMMIT:
S
Shengliang Guan 已提交
1276
        continueExec = mndTransPerformCommitStage(pMnode, pTrans);
S
Shengliang Guan 已提交
1277 1278
        break;
      case TRN_STAGE_ROLLBACK:
S
Shengliang Guan 已提交
1279 1280 1281 1282
        continueExec = mndTransPerformRollbackStage(pMnode, pTrans);
        break;
      case TRN_STAGE_FINISHED:
        continueExec = mndTransPerfromFinishedStage(pMnode, pTrans);
S
Shengliang Guan 已提交
1283
        break;
S
Shengliang Guan 已提交
1284
      default:
S
Shengliang Guan 已提交
1285 1286
        continueExec = false;
        break;
S
Shengliang Guan 已提交
1287 1288 1289
    }
  }

1290
  mndTransSendRpcRsp(pMnode, pTrans);
S
Shengliang Guan 已提交
1291
}
S
Shengliang Guan 已提交
1292

S
Shengliang Guan 已提交
1293 1294
static int32_t mndProcessTransReq(SRpcMsg *pReq) {
  mndTransPullup(pReq->info.node);
S
Shengliang Guan 已提交
1295 1296 1297
  return 0;
}

1298
int32_t mndKillTrans(SMnode *pMnode, STrans *pTrans) {
S
Shengliang Guan 已提交
1299 1300 1301 1302 1303 1304 1305 1306 1307 1308 1309 1310 1311 1312 1313
  SArray *pArray = NULL;
  if (pTrans->stage == TRN_STAGE_REDO_ACTION) {
    pArray = pTrans->redoActions;
  } else if (pTrans->stage == TRN_STAGE_UNDO_ACTION) {
    pArray = pTrans->undoActions;
  } else {
    terrno = TSDB_CODE_MND_TRANS_INVALID_STAGE;
    return -1;
  }

  int32_t size = taosArrayGetSize(pArray);

  for (int32_t i = 0; i < size; ++i) {
    STransAction *pAction = taosArrayGet(pArray, i);
    if (pAction == NULL) continue;
S
Shengliang Guan 已提交
1314

S
Shengliang Guan 已提交
1315
    if (pAction->msgReceived == 0) {
1316
      mInfo("trans:%d, action:%d set processed for kill msg received", pTrans->id, i);
S
Shengliang Guan 已提交
1317 1318 1319 1320 1321 1322
      pAction->msgSent = 1;
      pAction->msgReceived = 1;
      pAction->errCode = 0;
    }

    if (pAction->errCode != 0) {
1323
      mInfo("trans:%d, action:%d set processed for kill msg received, errCode from %s to success", pTrans->id, i,
S
Shengliang Guan 已提交
1324 1325 1326 1327 1328 1329 1330 1331 1332 1333 1334
            tstrerror(pAction->errCode));
      pAction->msgSent = 1;
      pAction->msgReceived = 1;
      pAction->errCode = 0;
    }
  }

  mndTransExecute(pMnode, pTrans);
  return 0;
}

S
Shengliang Guan 已提交
1335 1336
static int32_t mndProcessKillTransReq(SRpcMsg *pReq) {
  SMnode       *pMnode = pReq->info.node;
S
Shengliang Guan 已提交
1337
  SKillTransReq killReq = {0};
S
Shengliang Guan 已提交
1338
  int32_t       code = -1;
1339 1340
  SUserObj     *pUser = NULL;
  STrans       *pTrans = NULL;
S
Shengliang Guan 已提交
1341

S
Shengliang Guan 已提交
1342
  if (tDeserializeSKillTransReq(pReq->pCont, pReq->contLen, &killReq) != 0) {
S
Shengliang Guan 已提交
1343
    terrno = TSDB_CODE_INVALID_MSG;
S
Shengliang Guan 已提交
1344
    goto _OVER;
S
Shengliang Guan 已提交
1345 1346 1347 1348
  }

  mInfo("trans:%d, start to kill", killReq.transId);

S
Shengliang Guan 已提交
1349
  pUser = mndAcquireUser(pMnode, pReq->conn.user);
S
Shengliang Guan 已提交
1350
  if (pUser == NULL) {
S
Shengliang Guan 已提交
1351
    goto _OVER;
S
Shengliang Guan 已提交
1352 1353
  }

S
Shengliang Guan 已提交
1354
  if (mndCheckTransAuth(pUser) != 0) {
S
Shengliang Guan 已提交
1355
    goto _OVER;
S
Shengliang Guan 已提交
1356 1357
  }

S
Shengliang Guan 已提交
1358
  pTrans = mndAcquireTrans(pMnode, killReq.transId);
S
Shengliang Guan 已提交
1359 1360 1361 1362 1363 1364
  if (pTrans == NULL) {
    terrno = TSDB_CODE_MND_TRANS_NOT_EXIST;
    mError("trans:%d, failed to kill since %s", killReq.transId, terrstr());
    return -1;
  }

S
Shengliang Guan 已提交
1365 1366
  code = mndKillTrans(pMnode, pTrans);

S
Shengliang Guan 已提交
1367
_OVER:
1368
  if (code != 0) {
S
Shengliang Guan 已提交
1369 1370 1371 1372
    mError("trans:%d, failed to kill since %s", killReq.transId, terrstr());
    return -1;
  }

S
Shengliang Guan 已提交
1373
  mndReleaseTrans(pMnode, pTrans);
S
Shengliang Guan 已提交
1374
  return code;
S
Shengliang Guan 已提交
1375 1376
}

1377 1378
static int32_t mndCompareTransId(int32_t *pTransId1, int32_t *pTransId2) { return *pTransId1 >= *pTransId2 ? 1 : 0; }

S
Shengliang Guan 已提交
1379
void mndTransPullup(SMnode *pMnode) {
1380 1381 1382
  SSdb   *pSdb = pMnode->pSdb;
  SArray *pArray = taosArrayInit(sdbGetSize(pSdb, SDB_TRANS), sizeof(int32_t));
  if (pArray == NULL) return;
S
Shengliang Guan 已提交
1383

1384
  void *pIter = NULL;
S
Shengliang Guan 已提交
1385
  while (1) {
1386
    STrans *pTrans = NULL;
S
Shengliang Guan 已提交
1387 1388
    pIter = sdbFetch(pMnode->pSdb, SDB_TRANS, pIter, (void **)&pTrans);
    if (pIter == NULL) break;
1389 1390 1391
    taosArrayPush(pArray, &pTrans->id);
    sdbRelease(pSdb, pTrans);
  }
S
Shengliang Guan 已提交
1392

1393 1394 1395 1396 1397 1398 1399 1400 1401
  taosArraySort(pArray, (__compar_fn_t)mndCompareTransId);

  for (int32_t i = 0; i < taosArrayGetSize(pArray); ++i) {
    int32_t *pTransId = taosArrayGet(pArray, i);
    STrans  *pTrans = mndAcquireTrans(pMnode, *pTransId);
    if (pTrans != NULL) {
      mndTransExecute(pMnode, pTrans);
    }
    mndReleaseTrans(pMnode, pTrans);
S
Shengliang Guan 已提交
1402
  }
S
Shengliang Guan 已提交
1403 1404

  sdbWriteFile(pMnode->pSdb);
1405
  taosArrayDestroy(pArray);
1406
}
S
Shengliang Guan 已提交
1407

S
Shengliang Guan 已提交
1408 1409
static int32_t mndRetrieveTrans(SRpcMsg *pReq, SShowObj *pShow, SSDataBlock *pBlock, int32_t rows) {
  SMnode *pMnode = pReq->info.node;
1410
  SSdb   *pSdb = pMnode->pSdb;
S
Shengliang Guan 已提交
1411 1412 1413
  int32_t numOfRows = 0;
  STrans *pTrans = NULL;
  int32_t cols = 0;
1414
  char   *pWrite;
S
Shengliang Guan 已提交
1415 1416 1417 1418 1419 1420 1421

  while (numOfRows < rows) {
    pShow->pIter = sdbFetch(pSdb, SDB_TRANS, pShow->pIter, (void **)&pTrans);
    if (pShow->pIter == NULL) break;

    cols = 0;

S
Shengliang Guan 已提交
1422 1423
    SColumnInfoData *pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
    colDataAppend(pColInfo, numOfRows, (const char *)&pTrans->id, false);
S
Shengliang Guan 已提交
1424

S
Shengliang Guan 已提交
1425 1426
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
    colDataAppend(pColInfo, numOfRows, (const char *)&pTrans->createdTime, false);
S
Shengliang Guan 已提交
1427

S
Shengliang Guan 已提交
1428
    char stage[TSDB_TRANS_STAGE_LEN + VARSTR_HEADER_SIZE] = {0};
1429
    STR_WITH_MAXSIZE_TO_VARSTR(stage, mndTransStr(pTrans->stage), pShow->pMeta->pSchemas[cols].bytes);
S
Shengliang Guan 已提交
1430 1431
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
    colDataAppend(pColInfo, numOfRows, (const char *)stage, false);
S
Shengliang Guan 已提交
1432

S
Shengliang Guan 已提交
1433
    char dbname[TSDB_DB_NAME_LEN + VARSTR_HEADER_SIZE] = {0};
1434
    STR_WITH_MAXSIZE_TO_VARSTR(dbname, mndGetDbStr(pTrans->dbname), pShow->pMeta->pSchemas[cols].bytes);
S
Shengliang Guan 已提交
1435 1436
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
    colDataAppend(pColInfo, numOfRows, (const char *)dbname, false);
S
Shengliang Guan 已提交
1437

1438
    char type[TSDB_TRANS_TYPE_LEN + VARSTR_HEADER_SIZE] = {0};
1439
    STR_WITH_MAXSIZE_TO_VARSTR(type, mndTransType(pTrans->type), pShow->pMeta->pSchemas[cols].bytes);
S
Shengliang Guan 已提交
1440
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
1441
    colDataAppend(pColInfo, numOfRows, (const char *)type, false);
S
Shengliang Guan 已提交
1442

1443 1444 1445
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
    colDataAppend(pColInfo, numOfRows, (const char *)&pTrans->failedTimes, false);

S
Shengliang Guan 已提交
1446 1447
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
    colDataAppend(pColInfo, numOfRows, (const char *)&pTrans->lastExecTime, false);
S
Shengliang Guan 已提交
1448

S
Shengliang Guan 已提交
1449
    char lastError[TSDB_TRANS_ERROR_LEN + VARSTR_HEADER_SIZE] = {0};
1450
    STR_WITH_MAXSIZE_TO_VARSTR(lastError, pTrans->lastError, pShow->pMeta->pSchemas[cols].bytes);
S
Shengliang Guan 已提交
1451 1452
    pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
    colDataAppend(pColInfo, numOfRows, (const char *)lastError, false);
S
Shengliang Guan 已提交
1453 1454 1455 1456 1457

    numOfRows++;
    sdbRelease(pSdb, pTrans);
  }

1458
  pShow->numOfRows += numOfRows;
S
Shengliang Guan 已提交
1459 1460 1461 1462 1463 1464 1465
  return numOfRows;
}

static void mndCancelGetNextTrans(SMnode *pMnode, void *pIter) {
  SSdb *pSdb = pMnode->pSdb;
  sdbCancelFetch(pSdb, pIter);
}