mnodeSdb.c 28.0 KB
Newer Older
H
hzcheng 已提交
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15
/*
 * 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/>.
 */

S
slguan 已提交
16
#define _DEFAULT_SOURCE
17
#include "os.h"
S
slguan 已提交
18
#include "taoserror.h"
19
#include "hash.h"
S
slguan 已提交
20 21
#include "tutil.h"
#include "tbalance.h"
S
slguan 已提交
22
#include "tqueue.h"
S
slguan 已提交
23
#include "twal.h"
S
slguan 已提交
24
#include "tsync.h"
S
slguan 已提交
25
#include "tglobal.h"
S
slguan 已提交
26
#include "dnode.h"
27
#include "mnode.h"
S
Shengliang Guan 已提交
28 29 30 31 32
#include "mnodeDef.h"
#include "mnodeInt.h"
#include "mnodeMnode.h"
#include "mnodeDnode.h"
#include "mnodeSdb.h"
H
hzcheng 已提交
33

34
#define SDB_TABLE_LEN 12
35
#define SDB_SYNC_HACK 16
36

S
slguan 已提交
37 38 39 40 41 42 43 44 45
typedef enum {
  SDB_ACTION_INSERT,
  SDB_ACTION_DELETE,
  SDB_ACTION_UPDATE
} ESdbAction;

typedef enum {
  SDB_STATUS_OFFLINE,
  SDB_STATUS_SERVING,
S
slguan 已提交
46
  SDB_STATUS_CLOSING
S
slguan 已提交
47 48
} ESdbStatus;

S
slguan 已提交
49
typedef struct _SSdbTable {
50
  char      tableName[SDB_TABLE_LEN];
S
slguan 已提交
51 52 53 54 55 56 57 58 59 60 61 62 63 64 65
  ESdbTable tableId;
  ESdbKey   keyType;
  int32_t   hashSessions;
  int32_t   maxRowSize;
  int32_t   refCountPos;
  int32_t   autoIndex;
  int64_t   numOfRows;
  void *    iHandle;
  int32_t (*insertFp)(SSdbOper *pDesc);
  int32_t (*deleteFp)(SSdbOper *pOper);
  int32_t (*updateFp)(SSdbOper *pOper);
  int32_t (*decodeFp)(SSdbOper *pOper);
  int32_t (*encodeFp)(SSdbOper *pOper);
  int32_t (*destroyFp)(SSdbOper *pOper);
  int32_t (*restoredFp)();
S
slguan 已提交
66 67
} SSdbTable;

S
slguan 已提交
68 69 70 71 72 73 74 75 76 77 78 79 80 81
typedef struct {
  ESyncRole  role;
  ESdbStatus status;
  int64_t    version;
  void *     sync;
  void *     wal;
  SSyncCfg   cfg;
  sem_t      sem;
  int32_t    code;
  int32_t    numOfTables;
  SSdbTable *tableList[SDB_TABLE_MAX];
  pthread_mutex_t mutex;
} SSdbObject;

S
slguan 已提交
82
typedef struct {
S
slguan 已提交
83
  int32_t rowSize;
S
slguan 已提交
84
  void *  row;
S
slguan 已提交
85
} SSdbRow;
H
hzcheng 已提交
86

87 88 89 90 91 92 93 94 95 96
typedef struct {
  pthread_t thread;
  int32_t   workerId;
} SSdbWriteWorker;

typedef struct {
  int32_t num;
  SSdbWriteWorker *writeWorker;
} SSdbWriteWorkerPool;

S
slguan 已提交
97
static SSdbObject tsSdbObj = {0};
98 99 100 101 102 103 104 105 106 107 108 109
static taos_qset  tsSdbWriteQset;
static taos_qall  tsSdbWriteQall;
static taos_queue tsSdbWriteQueue;
static SSdbWriteWorkerPool tsSdbPool;

static int     sdbWrite(void *param, void *data, int type);
static int     sdbWriteToQueue(void *param, void *data, int type);
static void *  sdbWorkerFp(void *param);
static int32_t sdbInitWriteWorker();
static void    sdbCleanupWriteWorker();
static int32_t sdbAllocWriteQueue();
static void    sdbFreeWritequeue();
S
slguan 已提交
110

S
slguan 已提交
111 112 113 114 115 116 117
int32_t sdbGetId(void *handle) {
  return ((SSdbTable *)handle)->autoIndex;
}

int64_t sdbGetNumOfRows(void *handle) {
  return ((SSdbTable *)handle)->numOfRows;
}
S
slguan 已提交
118

S
slguan 已提交
119
uint64_t sdbGetVersion() {
S
slguan 已提交
120 121 122 123 124
  return tsSdbObj.version;
}

bool sdbIsMaster() { 
  return tsSdbObj.role == TAOS_SYNC_ROLE_MASTER; 
S
slguan 已提交
125 126
}

S
slguan 已提交
127 128 129 130
bool sdbIsServing() {
  return tsSdbObj.status == SDB_STATUS_SERVING; 
}

131 132 133 134 135 136 137 138
static void *sdbGetObjKey(SSdbTable *pTable, void *key) {
  if (pTable->keyType == SDB_KEY_VAR_STRING) {
    return *(char **)key;
  }

  return key;
}

S
slguan 已提交
139 140 141 142 143 144 145 146 147 148 149
static char *sdbGetActionStr(int32_t action) {
  switch (action) {
    case SDB_ACTION_INSERT:
      return "insert";
    case SDB_ACTION_DELETE:
      return "delete";
    case SDB_ACTION_UPDATE:
      return "update";
  }
  return "invalid";
}
S
slguan 已提交
150

151
static char *sdbGetKeyStr(SSdbTable *pTable, void *key) {
S
slguan 已提交
152 153
  static char str[16];
  switch (pTable->keyType) {
S
slguan 已提交
154
    case SDB_KEY_STRING:
155 156
    case SDB_KEY_VAR_STRING:
      return (char *)key;
S
slguan 已提交
157
    case SDB_KEY_INT:
S
slguan 已提交
158
    case SDB_KEY_AUTO:
159
      sprintf(str, "%d", *(int32_t *)key);
S
slguan 已提交
160 161
      return str;
    default:
S
slguan 已提交
162
      return "invalid";
S
slguan 已提交
163 164 165
  }
}

166 167 168 169
static char *sdbGetKeyStrFromObj(SSdbTable *pTable, void *key) {
  return sdbGetKeyStr(pTable, sdbGetObjKey(pTable, key));
}

S
slguan 已提交
170
static void *sdbGetTableFromId(int32_t tableId) {
S
slguan 已提交
171
  return tsSdbObj.tableList[tableId];
S
slguan 已提交
172 173
}

S
slguan 已提交
174
static int32_t sdbInitWal() {
H
hjxilinx 已提交
175
  SWalCfg walCfg = {.walLevel = 2, .wals = 2, .keep = 1};
S
slguan 已提交
176 177 178
  char temp[TSDB_FILENAME_LEN];
  sprintf(temp, "%s/wal", tsMnodeDir);
  tsSdbObj.wal = walOpen(temp, &walCfg);
S
slguan 已提交
179 180
  if (tsSdbObj.wal == NULL) {
    sdbError("failed to open sdb wal in %s", tsMnodeDir);
H
hzcheng 已提交
181 182 183
    return -1;
  }

S
slguan 已提交
184
  sdbTrace("open sdb wal for restore");
S
slguan 已提交
185
  walRestore(tsSdbObj.wal, NULL, sdbWrite);
S
slguan 已提交
186 187
  return 0;
}
H
hzcheng 已提交
188

S
slguan 已提交
189
static void sdbRestoreTables() {
S
slguan 已提交
190 191
  int32_t totalRows = 0;
  int32_t numOfTables = 0;
S
slguan 已提交
192
  for (int32_t tableId = 0; tableId < SDB_TABLE_MAX; ++tableId) {
S
slguan 已提交
193 194
    SSdbTable *pTable = sdbGetTableFromId(tableId);
    if (pTable == NULL) continue;
S
slguan 已提交
195 196
    if (pTable->restoredFp) {
      (*pTable->restoredFp)();
H
hzcheng 已提交
197 198
    }

S
slguan 已提交
199 200
    totalRows += pTable->numOfRows;
    numOfTables++;
201
    sdbTrace("table:%s, is restored, numOfRows:%" PRId64, pTable->tableName, pTable->numOfRows);
S
slguan 已提交
202 203
  }

204
  sdbTrace("sdb is restored, version:%" PRId64 " totalRows:%d numOfTables:%d", tsSdbObj.version, totalRows, numOfTables);
S
slguan 已提交
205 206 207 208 209 210 211 212
}

void sdbUpdateMnodeRoles() {
  if (tsSdbObj.sync == NULL) return;

  SNodesRole roles = {0};
  syncGetNodesRole(tsSdbObj.sync, &roles);

213
  sdbPrint("update mnodes sync roles, total:%d", tsSdbObj.cfg.replica);
S
slguan 已提交
214
  for (int32_t i = 0; i < tsSdbObj.cfg.replica; ++i) {
215
    SMnodeObj *pMnode = mnodeGetMnode(roles.nodeId[i]);
S
slguan 已提交
216 217
    if (pMnode != NULL) {
      pMnode->role = roles.role[i];
218
      sdbPrint("mnode:%d, role:%s", pMnode->mnodeId, mnodeGetMnodeRoleStr(pMnode->role));
219
      if (pMnode->mnodeId == dnodeGetDnodeId()) tsSdbObj.role = pMnode->role;
220
      mnodeDecMnodeRef(pMnode);
S
slguan 已提交
221 222
    }
  }
223

224
  mnodeUpdateMnodeIpSet();
S
slguan 已提交
225 226
}

陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
227
static uint32_t sdbGetFileInfo(void *ahandle, char *name, uint32_t *index, uint32_t eindex, int32_t *size, uint64_t *fversion) {
S
slguan 已提交
228 229 230 231 232
  sdbUpdateMnodeRoles();
  return 0;
}

static int sdbGetWalInfo(void *ahandle, char *name, uint32_t *index) {
S
slguan 已提交
233
  return walGetWalFile(tsSdbObj.wal, name, index);
S
slguan 已提交
234 235 236
}

static void sdbNotifyRole(void *ahandle, int8_t role) {
237
  sdbPrint("mnode role changed from %s to %s", mnodeGetMnodeRoleStr(tsSdbObj.role), mnodeGetMnodeRoleStr(role));
S
slguan 已提交
238 239 240 241 242 243 244 245 246 247 248 249

  if (role == TAOS_SYNC_ROLE_MASTER && tsSdbObj.role != TAOS_SYNC_ROLE_MASTER) {
    balanceReset();
  }
  tsSdbObj.role = role;

  sdbUpdateMnodeRoles();
}

static void sdbConfirmForward(void *ahandle, void *param, int32_t code) {
  tsSdbObj.code = code;
  sem_post(&tsSdbObj.sem);
S
slguan 已提交
250
  sdbTrace("forward request confirmed, version:%" PRIu64 ", result:%s", (int64_t)param, tstrerror(code));
S
slguan 已提交
251 252
}

S
Shengliang Guan 已提交
253
 static int32_t sdbForwardToPeer(SWalHead *pHead) {
S
slguan 已提交
254 255
  if (tsSdbObj.sync == NULL) return TSDB_CODE_SUCCESS;

256
  int32_t code = syncForwardToPeer(tsSdbObj.sync, pHead, (void*)pHead->version, TAOS_QTYPE_RPC);
S
slguan 已提交
257
  if (code > 0) {
258
    sdbTrace("forward request is sent, version:%" PRIu64 ", code:%d", pHead->version, code);
S
slguan 已提交
259 260 261 262 263 264 265 266 267 268
    sem_wait(&tsSdbObj.sem);
    return tsSdbObj.code;
  } 
  return code;
}

void sdbUpdateSync() {
  SSyncCfg syncCfg = {0};
  int32_t index = 0;

S
slguan 已提交
269
  SDMMnodeInfos *mnodes = dnodeGetMnodeInfos();
S
slguan 已提交
270
  for (int32_t i = 0; i < mnodes->nodeNum; ++i) {
S
slguan 已提交
271
    SDMMnodeInfo *node = &mnodes->nodeInfos[i];
S
slguan 已提交
272
    syncCfg.nodeInfo[i].nodeId = node->nodeId;
J
jtao1735 已提交
273 274
    taosGetFqdnPortFromEp(node->nodeEp, syncCfg.nodeInfo[i].nodeFqdn, &syncCfg.nodeInfo[i].nodePort);
    syncCfg.nodeInfo[i].nodePort += TSDB_PORT_SYNC;
S
slguan 已提交
275 276 277 278
    index++;
  }

  if (index == 0) {
S
Shengliang Guan 已提交
279
    void *pIter = NULL;
S
slguan 已提交
280 281
    while (1) {
      SMnodeObj *pMnode = NULL;
282
      pIter = mnodeGetNextMnode(pIter, &pMnode);
S
slguan 已提交
283 284 285 286
      if (pMnode == NULL) break;

      syncCfg.nodeInfo[index].nodeId = pMnode->mnodeId;

287
      SDnodeObj *pDnode = mnodeGetDnode(pMnode->mnodeId);
288 289 290 291 292 293
      if (pDnode != NULL) {
        syncCfg.nodeInfo[index].nodePort = pDnode->dnodePort + TSDB_PORT_SYNC;
        strcpy(syncCfg.nodeInfo[index].nodeFqdn, pDnode->dnodeEp);
        index++;
      }

294 295
      mnodeDecDnodeRef(pDnode);
      mnodeDecMnodeRef(pMnode);
S
slguan 已提交
296
    }
S
Shengliang Guan 已提交
297
    sdbFreeIter(pIter);
S
slguan 已提交
298 299 300
  }

  syncCfg.replica = index;
J
jtao1735 已提交
301
  syncCfg.quorum = (syncCfg.replica == 1) ? 1:2;
S
slguan 已提交
302 303 304 305 306 307 308 309 310 311 312 313

  bool hasThisDnode = false;
  for (int32_t i = 0; i < syncCfg.replica; ++i) {
    if (syncCfg.nodeInfo[i].nodeId == dnodeGetDnodeId()) {
      hasThisDnode = true;
      break;
    }
  }

  if (!hasThisDnode) return;
  if (memcmp(&syncCfg, &tsSdbObj.cfg, sizeof(SSyncCfg)) == 0) return;

J
jtao1735 已提交
314
  sdbPrint("work as mnode, replica:%d", syncCfg.replica);
S
slguan 已提交
315
  for (int32_t i = 0; i < syncCfg.replica; ++i) {
J
jtao1735 已提交
316
    sdbPrint("mnode:%d, %s:%d", syncCfg.nodeInfo[i].nodeId, syncCfg.nodeInfo[i].nodeFqdn, syncCfg.nodeInfo[i].nodePort);
S
slguan 已提交
317
  }
S
slguan 已提交
318

guanshengliang's avatar
guanshengliang 已提交
319
  SSyncInfo syncInfo = {0};
S
slguan 已提交
320 321 322
  syncInfo.vgId = 1;
  syncInfo.version = sdbGetVersion();
  syncInfo.syncCfg = syncCfg;
S
slguan 已提交
323
  sprintf(syncInfo.path, "%s", tsMnodeDir);
S
slguan 已提交
324 325 326
  syncInfo.ahandle = NULL;
  syncInfo.getWalInfo = sdbGetWalInfo;
  syncInfo.getFileInfo = sdbGetFileInfo;
327
  syncInfo.writeToCache = sdbWriteToQueue;
S
slguan 已提交
328 329 330
  syncInfo.confirmForward = sdbConfirmForward; 
  syncInfo.notifyRole = sdbNotifyRole;
  tsSdbObj.cfg = syncCfg;
331
  
S
slguan 已提交
332 333 334 335 336
  if (tsSdbObj.sync) {
    syncReconfig(tsSdbObj.sync, &syncCfg);
  } else {
    tsSdbObj.sync = syncStart(&syncInfo);
  }
337
  sdbUpdateMnodeRoles();
S
slguan 已提交
338 339 340 341 342 343
}

int32_t sdbInit() {
  pthread_mutex_init(&tsSdbObj.mutex, NULL);
  sem_init(&tsSdbObj.sem, 0, 0);

344 345 346 347
  if (sdbInitWriteWorker() != 0) {
    return -1;
  }

S
slguan 已提交
348 349 350
  if (sdbInitWal() != 0) {
    return -1;
  }
351

S
slguan 已提交
352
  sdbRestoreTables();
S
slguan 已提交
353

354
  if (mnodeGetMnodesNum() == 1) {
S
slguan 已提交
355 356 357 358 359 360
    tsSdbObj.role = TAOS_SYNC_ROLE_MASTER;
  }

  sdbUpdateSync();

  tsSdbObj.status = SDB_STATUS_SERVING;
S
slguan 已提交
361
  return TSDB_CODE_SUCCESS;
H
hzcheng 已提交
362 363
}

S
slguan 已提交
364
void sdbCleanUp() {
S
slguan 已提交
365 366
  if (tsSdbObj.status != SDB_STATUS_SERVING) return;

S
slguan 已提交
367
  tsSdbObj.status = SDB_STATUS_CLOSING;
guanshengliang's avatar
guanshengliang 已提交
368
  
369 370
  sdbCleanupWriteWorker();

guanshengliang's avatar
guanshengliang 已提交
371 372 373 374 375 376 377 378 379 380
  if (tsSdbObj.sync) {
    syncStop(tsSdbObj.sync);
    tsSdbObj.sync = NULL;
  }

  if (tsSdbObj.wal) {
    walClose(tsSdbObj.wal);
    tsSdbObj.wal = NULL;
  }
  
S
slguan 已提交
381 382
  sem_destroy(&tsSdbObj.sem);
  pthread_mutex_destroy(&tsSdbObj.mutex);
H
hzcheng 已提交
383 384
}

385 386 387 388 389 390
void sdbIncRef(void *handle, void *pObj) {
  if (pObj == NULL) return;

  SSdbTable *pTable = handle;
  int32_t *  pRefCount = (int32_t *)(pObj + pTable->refCountPos);
  atomic_add_fetch_32(pRefCount, 1);
391 392
  if (0 && (pTable->tableId == SDB_TABLE_CTABLE || pTable->tableId == SDB_TABLE_DB)) {
    sdbTrace("add ref to table:%s record:%p:%s:%d", pTable->tableName, pObj, sdbGetKeyStrFromObj(pTable, pObj), *pRefCount);
S
slguan 已提交
393 394 395
  }
}

396 397 398 399 400 401
void sdbDecRef(void *handle, void *pObj) {
  if (pObj == NULL) return;

  SSdbTable *pTable = handle;
  int32_t *  pRefCount = (int32_t *)(pObj + pTable->refCountPos);
  int32_t    refCount = atomic_sub_fetch_32(pRefCount, 1);
402 403
  if (0 && (pTable->tableId == SDB_TABLE_CTABLE || pTable->tableId == SDB_TABLE_DB)) {
    sdbTrace("def ref of table:%s record:%p:%s:%d", pTable->tableName, pObj, sdbGetKeyStrFromObj(pTable, pObj), *pRefCount);
S
slguan 已提交
404 405
  }

406 407
  int8_t *updateEnd = pObj + pTable->refCountPos - 1;
  if (refCount <= 0 && *updateEnd) {
408
    sdbTrace("table:%s, record:%p:%s:%d is destroyed", pTable->tableName, pObj, sdbGetKeyStrFromObj(pTable, pObj), *pRefCount);
409 410 411 412
    SSdbOper oper = {.pObj = pObj};
    (*pTable->destroyFp)(&oper);
  }
}
S
slguan 已提交
413

414 415
static SSdbRow *sdbGetRowMeta(SSdbTable *pTable, void *key) {
  if (pTable == NULL) return NULL;
S
slguan 已提交
416

417
  int32_t keySize = sizeof(int32_t);
418
  if (pTable->keyType == SDB_KEY_STRING || pTable->keyType == SDB_KEY_VAR_STRING) {
419 420
    keySize = strlen((char *)key);
  }
421 422 423
  
  return taosHashGet(pTable->iHandle, key, keySize);
}
S
slguan 已提交
424

425 426
static SSdbRow *sdbGetRowMetaFromObj(SSdbTable *pTable, void *key) {
  return sdbGetRowMeta(pTable, sdbGetObjKey(pTable, key));
S
slguan 已提交
427 428
}

H
hzcheng 已提交
429 430
void *sdbGetRow(void *handle, void *key) {
  SSdbTable *pTable = (SSdbTable *)handle;
431
  int32_t keySize = sizeof(int32_t);
432
  if (pTable->keyType == SDB_KEY_STRING || pTable->keyType == SDB_KEY_VAR_STRING) {
433 434
    keySize = strlen((char *)key);
  }
435 436 437 438 439 440 441 442
  
  SSdbRow *pMeta = taosHashGet(pTable->iHandle, key, keySize);
  if (pMeta) {
    sdbIncRef(pTable, pMeta->row);
    return pMeta->row;
  } else {
    return NULL;
  }
H
hzcheng 已提交
443 444
}

445 446 447 448
static void *sdbGetRowFromObj(SSdbTable *pTable, void *key) {
  return sdbGetRow(pTable, sdbGetObjKey(pTable, key));
}

S
slguan 已提交
449 450
static int32_t sdbInsertHash(SSdbTable *pTable, SSdbOper *pOper) {
  SSdbRow rowMeta;
S
slguan 已提交
451 452
  rowMeta.rowSize = pOper->rowSize;
  rowMeta.row = pOper->pObj;
H
hzcheng 已提交
453

454
  void *  key = sdbGetObjKey(pTable, pOper->pObj);
455
  int32_t keySize = sizeof(int32_t);
456 457 458

  if (pTable->keyType == SDB_KEY_STRING || pTable->keyType == SDB_KEY_VAR_STRING) {
    keySize = strlen((char *)key);
459
  }
460 461

  taosHashPut(pTable->iHandle, key, keySize, &rowMeta, sizeof(SSdbRow));
462

S
slguan 已提交
463
  sdbIncRef(pTable, pOper->pObj);
464
  atomic_add_fetch_32(&pTable->numOfRows, 1);
S
slguan 已提交
465

S
Shengliang Guan 已提交
466 467 468
  if (pTable->keyType == SDB_KEY_AUTO) {
    pTable->autoIndex = MAX(pTable->autoIndex, *((uint32_t *)pOper->pObj));
  } else {
469
    atomic_add_fetch_32(&pTable->autoIndex, 1);
S
slguan 已提交
470 471
  }

S
Shengliang Guan 已提交
472
  sdbTrace("table:%s, insert record:%s to hash, rowSize:%d numOfRows:%" PRId64 " version:%" PRIu64, pTable->tableName,
473
           sdbGetKeyStrFromObj(pTable, pOper->pObj), pOper->rowSize, pTable->numOfRows, sdbGetVersion());
S
slguan 已提交
474 475 476 477

  (*pTable->insertFp)(pOper);
  return TSDB_CODE_SUCCESS;
}
H
hzcheng 已提交
478

S
slguan 已提交
479
static int32_t sdbDeleteHash(SSdbTable *pTable, SSdbOper *pOper) {
S
slguan 已提交
480 481
  (*pTable->deleteFp)(pOper);
  
482
  void *  key = sdbGetObjKey(pTable, pOper->pObj);
483
  int32_t keySize = sizeof(int32_t);
484 485
  if (pTable->keyType == SDB_KEY_STRING || pTable->keyType == SDB_KEY_VAR_STRING) {
    keySize = strlen((char *)key);
486
  }
487 488

  taosHashRemove(pTable->iHandle, key, keySize);
489 490
  atomic_sub_fetch_32(&pTable->numOfRows, 1);
  
491
  sdbTrace("table:%s, delete record:%s from hash, numOfRows:%" PRId64 " version:%" PRIu64, pTable->tableName,
492
           sdbGetKeyStrFromObj(pTable, pOper->pObj), pTable->numOfRows, sdbGetVersion());
S
slguan 已提交
493

S
slguan 已提交
494
  int8_t *updateEnd = pOper->pObj + pTable->refCountPos - 1;
S
slguan 已提交
495 496
  *updateEnd = 1;
  sdbDecRef(pTable, pOper->pObj);
S
slguan 已提交
497

S
slguan 已提交
498 499 500
  return TSDB_CODE_SUCCESS;
}

S
slguan 已提交
501
static int32_t sdbUpdateHash(SSdbTable *pTable, SSdbOper *pOper) {
502
  sdbTrace("table:%s, update record:%s in hash, numOfRows:%" PRId64 " version:%" PRIu64, pTable->tableName,
503
           sdbGetKeyStrFromObj(pTable, pOper->pObj), pTable->numOfRows, sdbGetVersion());
S
slguan 已提交
504 505 506 507 508

  (*pTable->updateFp)(pOper);
  return TSDB_CODE_SUCCESS;
}

S
slguan 已提交
509
static int sdbWrite(void *param, void *data, int type) {
510
  SSdbOper *pOper = param;
S
slguan 已提交
511
  SWalHead *pHead = data;
512 513
  int32_t tableId = pHead->msgType / 10;
  int32_t action = pHead->msgType % 10;
S
slguan 已提交
514

S
slguan 已提交
515 516
  SSdbTable *pTable = sdbGetTableFromId(tableId);
  assert(pTable != NULL);
S
slguan 已提交
517

S
slguan 已提交
518 519 520 521 522 523 524 525 526
  pthread_mutex_lock(&tsSdbObj.mutex);
  if (pHead->version == 0) {
     // assign version
    tsSdbObj.version++;
    pHead->version = tsSdbObj.version;
  } else {
    // for data from WAL or forward, version may be smaller
    if (pHead->version <= tsSdbObj.version) {
      pthread_mutex_unlock(&tsSdbObj.mutex);
S
slguan 已提交
527 528 529 530
      if (type == TAOS_QTYPE_FWD && tsSdbObj.sync != NULL) {
        sdbTrace("forward request is received, version:%" PRIu64 " confirm it", pHead->version);
        syncConfirmForward(tsSdbObj.sync, pHead->version, TSDB_CODE_SUCCESS);
      }
S
slguan 已提交
531 532 533 534
      return TSDB_CODE_SUCCESS;
    } else if (pHead->version != tsSdbObj.version + 1) {
      pthread_mutex_unlock(&tsSdbObj.mutex);
      sdbError("table:%s, failed to restore %s record:%s from wal, version:%" PRId64 " too large, sdb version:%" PRId64,
535
               pTable->tableName, sdbGetActionStr(action), sdbGetKeyStr(pTable, pHead->cont), pHead->version,
S
slguan 已提交
536
               tsSdbObj.version);
537
      return TSDB_CODE_MND_APP_ERROR;
S
slguan 已提交
538 539 540
    } else {
      tsSdbObj.version = pHead->version;
    }
S
slguan 已提交
541
  }
S
slguan 已提交
542

S
slguan 已提交
543
  int32_t code = walWrite(tsSdbObj.wal, pHead);
S
slguan 已提交
544
  if (code < 0) {
S
slguan 已提交
545 546
    pthread_mutex_unlock(&tsSdbObj.mutex);
    return code;
S
slguan 已提交
547
  }
548
  
S
slguan 已提交
549
  code = sdbForwardToPeer(pHead);
S
slguan 已提交
550
  pthread_mutex_unlock(&tsSdbObj.mutex);
S
slguan 已提交
551

S
slguan 已提交
552
  // from app, oper is created
553
  if (pOper != NULL) {
554
    sdbTrace("record from app is disposed, version:%" PRIu64 " result:%s", pHead->version, tstrerror(code));
S
slguan 已提交
555 556 557 558
    return code;
  }
  
  // from wal or forward msg, oper not created, should add into hash
S
slguan 已提交
559
  if (tsSdbObj.sync != NULL) {
560
    sdbTrace("record from wal forward is disposed, version:%" PRIu64 " confirm it", pHead->version);
S
slguan 已提交
561
    syncConfirmForward(tsSdbObj.sync, pHead->version, code);
562 563
  } else {
    sdbTrace("record from wal restore is disposed, version:%" PRIu64 , pHead->version);
S
slguan 已提交
564
  }
S
slguan 已提交
565 566

  if (action == SDB_ACTION_INSERT) {
S
slguan 已提交
567
    SSdbOper oper = {.rowSize = pHead->len, .rowData = pHead->cont, .table = pTable};
S
slguan 已提交
568
    code = (*pTable->decodeFp)(&oper);
S
slguan 已提交
569
    return sdbInsertHash(pTable, &oper);
S
slguan 已提交
570
  } else if (action == SDB_ACTION_DELETE) {
S
slguan 已提交
571
    SSdbRow *rowMeta = sdbGetRowMeta(pTable, pHead->cont);
S
slguan 已提交
572
    assert(rowMeta != NULL && rowMeta->row != NULL);
S
slguan 已提交
573 574
    SSdbOper oper = {.table = pTable, .pObj = rowMeta->row};
    return sdbDeleteHash(pTable, &oper);
S
slguan 已提交
575
  } else if (action == SDB_ACTION_UPDATE) {
S
slguan 已提交
576
    SSdbRow *rowMeta = sdbGetRowMeta(pTable, pHead->cont);
S
slguan 已提交
577
    assert(rowMeta != NULL && rowMeta->row != NULL);
S
slguan 已提交
578
    SSdbOper oper = {.rowSize = pHead->len, .rowData = pHead->cont, .table = pTable};
S
slguan 已提交
579
    code = (*pTable->decodeFp)(&oper);
S
slguan 已提交
580
    return sdbUpdateHash(pTable, &oper);
581
  } else { return TSDB_CODE_MND_INVALID_MSG_TYPE; }
S
slguan 已提交
582 583
}

S
slguan 已提交
584
int32_t sdbInsertRow(SSdbOper *pOper) {
S
slguan 已提交
585
  SSdbTable *pTable = (SSdbTable *)pOper->table;
586
  if (pTable == NULL) return TSDB_CODE_MND_SDB_INVALID_TABLE_TYPE;
S
slguan 已提交
587

588 589
  if (sdbGetRowFromObj(pTable, pOper->pObj)) {
    sdbError("table:%s, failed to insert record:%s, already exist", pTable->tableName, sdbGetKeyStrFromObj(pTable, pOper->pObj));
S
slguan 已提交
590
    sdbDecRef(pTable, pOper->pObj);
591
    return TSDB_CODE_MND_SDB_OBJ_ALREADY_THERE;
S
slguan 已提交
592
  }
S
slguan 已提交
593

S
slguan 已提交
594
  if (pTable->keyType == SDB_KEY_AUTO) {
595
    *((uint32_t *)pOper->pObj) = atomic_add_fetch_32(&pTable->autoIndex, 1);
S
slguan 已提交
596 597 598

    // let vgId increase from 2
    if (pTable->autoIndex == 1 && strcmp(pTable->tableName, "vgroups") == 0) {
599
      *((uint32_t *)pOper->pObj) = atomic_add_fetch_32(&pTable->autoIndex, 1);
S
slguan 已提交
600
    }
S
slguan 已提交
601
  }
H
hzcheng 已提交
602

603 604 605 606 607
  int32_t code = sdbInsertHash(pTable, pOper);
  if (code != TSDB_CODE_SUCCESS) {
    sdbError("table:%s, failed to insert into hash", pTable->tableName);
    return code;
  }
S
slguan 已提交
608

609 610 611
  // just insert data into memory
  if (pOper->type != SDB_OPER_GLOBAL) {
    return TSDB_CODE_SUCCESS;
S
slguan 已提交
612 613
  }

614
  int32_t size = sizeof(SSdbOper) + sizeof(SWalHead) + pTable->maxRowSize + SDB_SYNC_HACK;
615 616
  SSdbOper *pNewOper = taosAllocateQitem(size);
  
617
  SWalHead *pHead = (void *)pNewOper + sizeof(SSdbOper) + SDB_SYNC_HACK;
618 619 620 621 622 623 624 625 626
  pHead->version = 0;
  pHead->len = pOper->rowSize;
  pHead->msgType = pTable->tableId * 10 + SDB_ACTION_INSERT;

  pOper->rowData = pHead->cont;
  (*pTable->encodeFp)(pOper);
  pHead->len = pOper->rowSize;

  memcpy(pNewOper, pOper, sizeof(SSdbOper));
627 628

  if (pNewOper->pMsg != NULL) {
S
Shengliang Guan 已提交
629
    sdbTrace("app:%p:%p, insert action is add to sdb queue", pNewOper->pMsg->rpcMsg.ahandle, pNewOper->pMsg);
630 631
  }

632 633
  taosWriteQitem(tsSdbWriteQueue, TAOS_QTYPE_RPC, pNewOper);
  return TSDB_CODE_SUCCESS;
H
hzcheng 已提交
634 635
}

S
slguan 已提交
636
int32_t sdbDeleteRow(SSdbOper *pOper) {
S
slguan 已提交
637
  SSdbTable *pTable = (SSdbTable *)pOper->table;
638
  if (pTable == NULL) return TSDB_CODE_MND_SDB_INVALID_TABLE_TYPE;
H
hzcheng 已提交
639

640
  SSdbRow *pMeta = sdbGetRowMetaFromObj(pTable, pOper->pObj);
H
hzcheng 已提交
641
  if (pMeta == NULL) {
S
slguan 已提交
642
    sdbTrace("table:%s, record is not there, delete failed", pTable->tableName);
643
    return TSDB_CODE_MND_SDB_OBJ_NOT_THERE;
H
hzcheng 已提交
644 645
  }

646 647 648 649 650
  void *pMetaRow = pMeta->row;
  if (pMetaRow == NULL) {
    sdbError("table:%s, record meta is null", pTable->tableName);
    return TSDB_CODE_MND_SDB_INVAID_META_ROW;
  }
S
slguan 已提交
651

652 653 654 655 656
  int32_t code = sdbDeleteHash(pTable, pOper);
  if (code != TSDB_CODE_SUCCESS) {
    sdbError("table:%s, failed to delete from hash", pTable->tableName);
    return code;
  }
H
hzcheng 已提交
657

658 659 660
  // just delete data from memory
  if (pOper->type != SDB_OPER_GLOBAL) {
    return TSDB_CODE_SUCCESS;
S
slguan 已提交
661 662
  }

663
  int32_t size = sizeof(SSdbOper) + sizeof(SWalHead) + pTable->maxRowSize + SDB_SYNC_HACK;
664 665
  SSdbOper *pNewOper = taosAllocateQitem(size);

666
  SWalHead *pHead = (void *)pNewOper + sizeof(SSdbOper) + SDB_SYNC_HACK;
667 668
  pHead->version = 0;
  pHead->msgType = pTable->tableId * 10 + SDB_ACTION_DELETE;
669 670 671 672
  
  pOper->rowData = pHead->cont;
  (*pTable->encodeFp)(pOper);
  pHead->len = pOper->rowSize;
673 674

  memcpy(pNewOper, pOper, sizeof(SSdbOper));
675 676

  if (pNewOper->pMsg != NULL) {
S
Shengliang Guan 已提交
677
    sdbTrace("app:%p:%p, delete action is add to sdb queue", pNewOper->pMsg->rpcMsg.ahandle, pNewOper->pMsg);
678 679
  }

680 681
  taosWriteQitem(tsSdbWriteQueue, TAOS_QTYPE_RPC, pNewOper);
  return TSDB_CODE_SUCCESS;
H
hzcheng 已提交
682 683
}

S
slguan 已提交
684
int32_t sdbUpdateRow(SSdbOper *pOper) {
S
slguan 已提交
685
  SSdbTable *pTable = (SSdbTable *)pOper->table;
686
  if (pTable == NULL) return TSDB_CODE_MND_SDB_INVALID_TABLE_TYPE;
H
hzcheng 已提交
687

688
  SSdbRow *pMeta = sdbGetRowMetaFromObj(pTable, pOper->pObj);
H
hzcheng 已提交
689
  if (pMeta == NULL) {
690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708
    sdbTrace("table:%s, record is not there, update failed", pTable->tableName);
    return TSDB_CODE_MND_SDB_OBJ_NOT_THERE;
  }

  void *pMetaRow = pMeta->row;
  if (pMetaRow == NULL) {
    sdbError("table:%s, record meta is null", pTable->tableName);
    return TSDB_CODE_MND_SDB_INVAID_META_ROW;
  }

  int32_t code = sdbUpdateHash(pTable, pOper);
  if (code != TSDB_CODE_SUCCESS) {
    sdbError("table:%s, failed to update hash", pTable->tableName);
    return code;
  }

  // just update data in memory
  if (pOper->type != SDB_OPER_GLOBAL) {
    return TSDB_CODE_SUCCESS;
H
hzcheng 已提交
709 710
  }

711
  int32_t size = sizeof(SSdbOper) + sizeof(SWalHead) + pTable->maxRowSize + SDB_SYNC_HACK;
712
  SSdbOper *pNewOper = taosAllocateQitem(size);
H
hzcheng 已提交
713

714
  SWalHead *pHead = (void *)pNewOper + sizeof(SSdbOper) + SDB_SYNC_HACK;
715 716
  pHead->version = 0;
  pHead->msgType = pTable->tableId * 10 + SDB_ACTION_UPDATE;
H
hzcheng 已提交
717

718 719 720
  pOper->rowData = pHead->cont;
  (*pTable->encodeFp)(pOper);
  pHead->len = pOper->rowSize;
S
slguan 已提交
721

722 723 724
  memcpy(pNewOper, pOper, sizeof(SSdbOper));

  if (pNewOper->pMsg != NULL) {
S
Shengliang Guan 已提交
725
    sdbTrace("app:%p:%p, update action is add to sdb queue", pNewOper->pMsg->rpcMsg.ahandle, pNewOper->pMsg);
726 727
  }

728 729
  taosWriteQitem(tsSdbWriteQueue, TAOS_QTYPE_RPC, pNewOper);
  return TSDB_CODE_SUCCESS;
S
slguan 已提交
730
}
S
slguan 已提交
731

S
slguan 已提交
732 733 734 735
void *sdbFetchRow(void *handle, void *pNode, void **ppRow) {
  SSdbTable *pTable = (SSdbTable *)handle;
  *ppRow = NULL;
  if (pTable == NULL) return NULL;
H
hzcheng 已提交
736

737 738 739 740 741 742 743 744 745 746 747
  SHashMutableIterator *pIter = pNode;
  if (pIter == NULL) {
    pIter = taosHashCreateIter(pTable->iHandle);
  }

  if (!taosHashIterNext(pIter)) {
    taosHashDestroyIter(pIter);
    return NULL;
  }

  SSdbRow *pMeta = taosHashIterGet(pIter);
S
Shengliang Guan 已提交
748 749 750 751
  if (pMeta == NULL) {
    taosHashDestroyIter(pIter);
    return NULL;
  }
S
slguan 已提交
752

S
slguan 已提交
753 754
  *ppRow = pMeta->row;
  sdbIncRef(handle, pMeta->row);
H
hzcheng 已提交
755

756
  return pIter;
S
slguan 已提交
757 758
}

S
Shengliang Guan 已提交
759 760 761 762 763 764
void sdbFreeIter(void *pIter) {
  if (pIter != NULL) {
    taosHashDestroyIter(pIter);
  }
}

S
slguan 已提交
765 766
void *sdbOpenTable(SSdbTableDesc *pDesc) {
  SSdbTable *pTable = (SSdbTable *)calloc(1, sizeof(SSdbTable));
S
slguan 已提交
767
  
S
slguan 已提交
768 769
  if (pTable == NULL) return NULL;

770
  tstrncpy(pTable->tableName, pDesc->tableName, SDB_TABLE_LEN);
S
slguan 已提交
771 772 773 774 775 776 777 778 779 780 781
  pTable->keyType      = pDesc->keyType;
  pTable->tableId      = pDesc->tableId;
  pTable->hashSessions = pDesc->hashSessions;
  pTable->maxRowSize   = pDesc->maxRowSize;
  pTable->refCountPos  = pDesc->refCountPos;
  pTable->insertFp     = pDesc->insertFp;
  pTable->deleteFp     = pDesc->deleteFp;
  pTable->updateFp     = pDesc->updateFp;
  pTable->encodeFp     = pDesc->encodeFp;
  pTable->decodeFp     = pDesc->decodeFp;
  pTable->destroyFp    = pDesc->destroyFp;
S
slguan 已提交
782
  pTable->restoredFp   = pDesc->restoredFp;
783 784

  _hash_fn_t hashFp = taosGetDefaultHashFunction(TSDB_DATA_TYPE_INT);
S
Shengliang Guan 已提交
785
  if (pTable->keyType == SDB_KEY_STRING || pTable->keyType == SDB_KEY_VAR_STRING) {
786
    hashFp = taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY);
S
slguan 已提交
787
  }
788
  pTable->iHandle = taosHashInit(pTable->hashSessions, hashFp, true);
S
slguan 已提交
789

S
slguan 已提交
790 791
  tsSdbObj.numOfTables++;
  tsSdbObj.tableList[pTable->tableId] = pTable;
S
slguan 已提交
792
  return pTable;
H
hzcheng 已提交
793 794 795 796 797
}

void sdbCloseTable(void *handle) {
  SSdbTable *pTable = (SSdbTable *)handle;
  if (pTable == NULL) return;
S
slguan 已提交
798
  
S
slguan 已提交
799 800
  tsSdbObj.numOfTables--;
  tsSdbObj.tableList[pTable->tableId] = NULL;
H
hzcheng 已提交
801

802 803 804 805
  SHashMutableIterator *pIter = taosHashCreateIter(pTable->iHandle);
  while (taosHashIterNext(pIter)) {
    SSdbRow *pMeta = taosHashIterGet(pIter);
    if (pMeta == NULL) continue;
S
slguan 已提交
806

S
slguan 已提交
807
    SSdbOper oper = {
808
      .pObj = pMeta->row,
S
slguan 已提交
809 810
      .table = pTable,
    };
811

S
slguan 已提交
812
    (*pTable->destroyFp)(&oper);
H
hzcheng 已提交
813 814
  }

815 816
  taosHashDestroyIter(pIter);
  taosHashCleanup(pTable->iHandle);
H
hzcheng 已提交
817

S
slguan 已提交
818
  sdbTrace("table:%s, is closed, numOfTables:%d", pTable->tableName, tsSdbObj.numOfTables);
S
slguan 已提交
819
  free(pTable);
H
hzcheng 已提交
820
}
S
slguan 已提交
821

822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853
int32_t sdbInitWriteWorker() {
  tsSdbPool.num = 1;
  tsSdbPool.writeWorker = (SSdbWriteWorker *)calloc(sizeof(SSdbWriteWorker), tsSdbPool.num);

  if (tsSdbPool.writeWorker == NULL) return -1;
  for (int32_t i = 0; i < tsSdbPool.num; ++i) {
    SSdbWriteWorker *pWorker = tsSdbPool.writeWorker + i;
    pWorker->workerId = i;
  }

  sdbAllocWriteQueue();
  
  mPrint("sdb write is opened");
  return 0;
}

void sdbCleanupWriteWorker() {
  for (int32_t i = 0; i < tsSdbPool.num; ++i) {
    SSdbWriteWorker *pWorker = tsSdbPool.writeWorker + i;
    if (pWorker->thread) {
      taosQsetThreadResume(tsSdbWriteQset);
    }
  }

  for (int32_t i = 0; i < tsSdbPool.num; ++i) {
    SSdbWriteWorker *pWorker = tsSdbPool.writeWorker + i;
    if (pWorker->thread) {
      pthread_join(pWorker->thread, NULL);
    }
  }

  sdbFreeWritequeue();
陶建辉(Jeff)'s avatar
陶建辉(Jeff) 已提交
854
  tfree(tsSdbPool.writeWorker);
855 856 857 858 859 860 861 862 863 864 865 866 867 868 869 870 871 872 873 874 875 876 877 878 879 880 881 882 883 884 885 886 887 888 889 890 891 892 893 894 895 896 897 898 899 900 901 902 903 904 905 906 907 908 909 910 911 912 913 914 915 916 917 918 919 920 921 922 923 924 925 926 927 928 929 930 931 932 933 934 935 936 937 938

  mPrint("sdb write is closed");
}

int32_t sdbAllocWriteQueue() {
  tsSdbWriteQueue = taosOpenQueue();
  if (tsSdbWriteQueue == NULL) return TSDB_CODE_MND_OUT_OF_MEMORY;

  tsSdbWriteQset = taosOpenQset();
  if (tsSdbWriteQset == NULL) {
    taosCloseQueue(tsSdbWriteQueue);
    return TSDB_CODE_MND_OUT_OF_MEMORY;
  }
  taosAddIntoQset(tsSdbWriteQset, tsSdbWriteQueue, NULL);

  tsSdbWriteQall = taosAllocateQall();
  if (tsSdbWriteQall == NULL) {
    taosCloseQset(tsSdbWriteQset);
    taosCloseQueue(tsSdbWriteQueue);
    return TSDB_CODE_MND_OUT_OF_MEMORY;
  }
  
  for (int32_t i = 0; i < tsSdbPool.num; ++i) {
    SSdbWriteWorker *pWorker = tsSdbPool.writeWorker + i;
    pWorker->workerId = i;

    pthread_attr_t thAttr;
    pthread_attr_init(&thAttr);
    pthread_attr_setdetachstate(&thAttr, PTHREAD_CREATE_JOINABLE);

    if (pthread_create(&pWorker->thread, &thAttr, sdbWorkerFp, pWorker) != 0) {
      mError("failed to create thread to process sdb write queue, reason:%s", strerror(errno));
      taosFreeQall(tsSdbWriteQall);
      taosCloseQset(tsSdbWriteQset);
      taosCloseQueue(tsSdbWriteQueue);
      return TSDB_CODE_MND_OUT_OF_MEMORY;
    }

    pthread_attr_destroy(&thAttr);
    mTrace("sdb write worker:%d is launched, total:%d", pWorker->workerId, tsSdbPool.num);
  }

  mTrace("sdb write queue:%p is allocated", tsSdbWriteQueue);
  return TSDB_CODE_SUCCESS;
}

void sdbFreeWritequeue() {
  taosCloseQset(tsSdbWriteQueue);
  taosFreeQall(tsSdbWriteQall);
  taosCloseQset(tsSdbWriteQset);
  tsSdbWriteQall = NULL;
  tsSdbWriteQset = NULL;
  tsSdbWriteQueue = NULL;
}

int sdbWriteToQueue(void *param, void *data, int type) {
  SWalHead *pHead = data;
  int size = sizeof(SWalHead) + pHead->len;
  SWalHead *pWal = (SWalHead *)taosAllocateQitem(size);
  memcpy(pWal, pHead, size);

  taosWriteQitem(tsSdbWriteQueue, type, pWal);
  return 0;
}

static void *sdbWorkerFp(void *param) {
  SWalHead *pHead;
  SSdbOper *pOper;
  int32_t   type;
  int32_t   numOfMsgs;
  void *    item;
  void *    unUsed;

  while (1) {
    numOfMsgs = taosReadAllQitemsFromQset(tsSdbWriteQset, tsSdbWriteQall, &unUsed);
    if (numOfMsgs == 0) {
      sdbTrace("sdbWorkerFp: got no message from qset, exiting...");
      break;
    }

    for (int32_t i = 0; i < numOfMsgs; ++i) {
      taosGetQitem(tsSdbWriteQall, &type, &item);
      if (type == TAOS_QTYPE_RPC) {
        pOper = (SSdbOper *)item;
939
        pHead = (void *)pOper + sizeof(SSdbOper) + SDB_SYNC_HACK;
940 941 942 943 944
      } else {
        pHead = (SWalHead *)item;
        pOper = NULL;
      }

945 946 947 948
      if (pOper != NULL && pOper->pMsg != NULL) {
        sdbTrace("app:%p:%p, will be processed in sdb queue", pOper->pMsg->rpcMsg.ahandle, pOper->pMsg);
      }

949
      int32_t code = sdbWrite(pOper, pHead, type);
950
      if (pOper) pOper->retCode = code;
951 952 953 954 955 956 957 958 959 960
    }

    walFsync(tsSdbObj.wal);

    // browse all items, and process them one by one
    taosResetQitems(tsSdbWriteQall);
    for (int32_t i = 0; i < numOfMsgs; ++i) {
      taosGetQitem(tsSdbWriteQall, &type, &item);
      if (type == TAOS_QTYPE_RPC) {
        pOper = (SSdbOper *)item;
961
        if (pOper != NULL && pOper->cb != NULL) {
962
          pOper->retCode = (*pOper->cb)(pOper->pMsg, pOper->retCode);
963
        }
964
          
965 966 967 968 969
        if (pOper != NULL && pOper->pMsg != NULL) {
          sdbTrace("app:%p:%p, msg is processed, result:%s", pOper->pMsg->rpcMsg.ahandle, pOper->pMsg,
                   tstrerror(pOper->retCode));
        }

970
        dnodeSendRpcMnodeWriteRsp(pOper->pMsg, pOper->retCode);
971 972 973 974 975 976 977
      }
      taosFreeQitem(item);
    }
  }

  return NULL;
}