• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

taosdata / TDengine / #5089

17 May 2026 01:15AM UTC coverage: 73.286% (-0.05%) from 73.335%
#5089

push

travis-ci

web-flow
feat (TDgpt): Dynamic Model Synchronization Enhancements (#35344)

* refactor: do some internal refactor.

* fix: fix multiprocess sync issue.

* feat: add dynamic anomaly detection and forecasting services

* fix: log error message for undeploying model in exception handling

* Potential fix for pull request finding

Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>

* Potential fix for pull request finding

Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>

* Potential fix for pull request finding

Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>

* Potential fix for pull request finding

Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>

* fix: handle undeploy when model exists only on disk

Agent-Logs-Url: https://github.com/taosdata/TDengine/sessions/286aafa0-c3ce-4c27-b803-2707571e9dc1

Co-authored-by: hjxilinx <8252296+hjxilinx@users.noreply.github.com>

* fix: guard dynamic registry concurrent access

Agent-Logs-Url: https://github.com/taosdata/TDengine/sessions/5e4db858-6458-40f4-ac28-d1b1b7f97c18

Co-authored-by: hjxilinx <8252296+hjxilinx@users.noreply.github.com>

* fix: tighten service list locking scope

Agent-Logs-Url: https://github.com/taosdata/TDengine/sessions/5e4db858-6458-40f4-ac28-d1b1b7f97c18

Co-authored-by: hjxilinx <8252296+hjxilinx@users.noreply.github.com>

* fix: restore prophet support and update tests per review feedback

Agent-Logs-Url: https://github.com/taosdata/TDengine/sessions/92298ae1-7da6-4d07-b20e-101c7cd0b26b

Co-authored-by: hjxilinx <8252296+hjxilinx@users.noreply.github.com>

* fix: improve test name and move copy inside lock scope

Agent-Logs-Url: https://github.com/taosdata/TDengine/sessions/92298ae1-7da6-4d07-b20e-101c7cd0b26b

Co-authored-by: hjxilinx <8252296+hjxilinx@users.noreply.github.com>

* Potential fix for pull request finding

Co-au... (continued)

281267 of 383795 relevant lines covered (73.29%)

135515062.05 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

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

16
#include "mndTopic.h"
17
#include "audit.h"
18
#include "mndConsumer.h"
19
#include "mndDb.h"
20
#include "mndDnode.h"
21
#include "mndMnode.h"
22
#include "mndPrivilege.h"
23
#include "mndShow.h"
24
#include "mndStb.h"
25
#include "mndSubscribe.h"
26
#include "mndTrans.h"
27
#include "mndUser.h"
28
#include "mndVgroup.h"
29
#include "osMemPool.h"
30
#include "parser.h"
31
#include "tlockfree.h"
32
#include "tname.h"
33

34
#define MND_TOPIC_VER_SUPPORT_OWNER 4
35
#define MND_TOPIC_VER_NUMBER        MND_TOPIC_VER_SUPPORT_OWNER
36
#define MND_TOPIC_RESERVE_SIZE      64
37

38
SHashObj *topicsToReload = NULL;
39

40
SSdbRaw *mndTopicActionEncode(SMqTopicObj *pTopic);
41
SSdbRow *mndTopicActionDecode(SSdbRaw *pRaw);
42

43
static int32_t mndTopicActionInsert(SSdb *pSdb, SMqTopicObj *pTopic);
44
static int32_t mndTopicActionDelete(SSdb *pSdb, SMqTopicObj *pTopic);
45
static int32_t mndTopicActionUpdate(SSdb *pSdb, SMqTopicObj *pOldTopic, SMqTopicObj *pNewTopic);
46
static int32_t mndProcessCreateTopicReq(SRpcMsg *pReq);
47
static int32_t mndProcessDropTopicReq(SRpcMsg *pReq);
48

49
static int32_t mndRetrieveTopic(SRpcMsg *pReq, SShowObj *pShow, SSDataBlock *pBlock, int32_t rows);
50
static void    mndCancelGetNextTopic(SMnode *pMnode, void *pIter);
51
static int32_t processAst(SMqTopicObj *topicObj, const char *ast);
52

53
int32_t mndInitTopic(SMnode *pMnode) {
528,435✔
54
  SSdbTable table = {
528,435✔
55
      .sdbType = SDB_TOPIC,
56
      .keyType = SDB_KEY_BINARY,
57
      .encodeFp = (SdbEncodeFp)mndTopicActionEncode,
58
      .decodeFp = (SdbDecodeFp)mndTopicActionDecode,
59
      .insertFp = (SdbInsertFp)mndTopicActionInsert,
60
      .updateFp = (SdbUpdateFp)mndTopicActionUpdate,
61
      .deleteFp = (SdbDeleteFp)mndTopicActionDelete,
62
  };
63

64
  if (pMnode == NULL) {
528,435✔
65
    return TSDB_CODE_INVALID_PARA;
×
66
  }
67
  mndSetMsgHandle(pMnode, TDMT_MND_TMQ_CREATE_TOPIC, mndProcessCreateTopicReq);
528,435✔
68
  mndSetMsgHandle(pMnode, TDMT_MND_TMQ_DROP_TOPIC, mndProcessDropTopicReq);
528,435✔
69

70
  mndAddShowRetrieveHandle(pMnode, TSDB_MGMT_TABLE_TOPICS, mndRetrieveTopic);
528,435✔
71
  mndAddShowFreeIterHandle(pMnode, TSDB_MGMT_TABLE_TOPICS, mndCancelGetNextTopic);
528,435✔
72

73
  return sdbSetTable(pMnode->pSdb, table);
528,435✔
74
}
75

76
void mndCleanupTopic(SMnode *pMnode) {}
528,373✔
77

78
void mndTopicGetShowName(const char *fullTopic, char *topic) {
185,259✔
79
  if (fullTopic == NULL) {
185,259✔
80
    return;
×
81
  }
82
  char *tmp = strchr(fullTopic, '.');
185,259✔
83
  if (tmp == NULL) {
185,259✔
84
    tstrncpy(topic, fullTopic, TSDB_TOPIC_FNAME_LEN);
×
85
  } else {
86
    tstrncpy(topic, tmp + 1, TSDB_TOPIC_FNAME_LEN);
185,259✔
87
  }
88
}
89
SSdbRaw *mndTopicActionEncode(SMqTopicObj *pTopic) {
362,571✔
90
  if (pTopic == NULL) {
362,571✔
91
    return NULL;
×
92
  }
93
  int32_t code = 0;
362,571✔
94
  int32_t lino = 0;
362,571✔
95
  terrno = TSDB_CODE_OUT_OF_MEMORY;
362,571✔
96

97
  void *  swBuf = NULL;
362,571✔
98
  int32_t physicalPlanLen = 0;
362,571✔
99
  if (pTopic->physicalPlan) {
362,571✔
100
    physicalPlanLen = strlen(pTopic->physicalPlan) + 1;
362,571✔
101
  }
102

103
  int32_t schemaLen = 0;
362,571✔
104
  if (pTopic->schema.nCols) {
362,571✔
105
    schemaLen = taosEncodeSSchemaWrapper(NULL, &pTopic->schema);
555,938✔
106
  }
107

108
  int32_t  size = sizeof(SMqTopicObj) + physicalPlanLen + pTopic->sqlLen + schemaLen + MND_TOPIC_RESERVE_SIZE;
362,571✔
109
  SSdbRaw *pRaw = sdbAllocRaw(SDB_TOPIC, MND_TOPIC_VER_NUMBER, size);
362,571✔
110
  if (pRaw == NULL) {
362,571✔
111
    goto TOPIC_ENCODE_OVER;
×
112
  }
113

114
  int32_t dataPos = 0;
362,571✔
115
  SDB_SET_BINARY(pRaw, dataPos, pTopic->name, TSDB_TOPIC_FNAME_LEN, TOPIC_ENCODE_OVER);
362,571✔
116
  SDB_SET_BINARY(pRaw, dataPos, pTopic->db, TSDB_DB_FNAME_LEN, TOPIC_ENCODE_OVER);
362,571✔
117
  SDB_SET_BINARY(pRaw, dataPos, pTopic->createUser, TSDB_USER_LEN, TOPIC_ENCODE_OVER);
362,571✔
118
  SDB_SET_INT64(pRaw, dataPos, pTopic->createTime, TOPIC_ENCODE_OVER);
362,571✔
119
  SDB_SET_INT64(pRaw, dataPos, pTopic->updateTime, TOPIC_ENCODE_OVER);
362,571✔
120
  SDB_SET_INT64(pRaw, dataPos, pTopic->uid, TOPIC_ENCODE_OVER);
362,571✔
121
  SDB_SET_INT64(pRaw, dataPos, pTopic->dbUid, TOPIC_ENCODE_OVER);
362,571✔
122
  SDB_SET_INT32(pRaw, dataPos, pTopic->version, TOPIC_ENCODE_OVER);
362,571✔
123
  SDB_SET_INT8(pRaw, dataPos, pTopic->subType, TOPIC_ENCODE_OVER);
362,571✔
124
  SDB_SET_INT8(pRaw, dataPos, pTopic->withMeta, TOPIC_ENCODE_OVER);
362,571✔
125

126
  SDB_SET_INT64(pRaw, dataPos, pTopic->stbUid, TOPIC_ENCODE_OVER);
362,571✔
127
  SDB_SET_BINARY(pRaw, dataPos, pTopic->stbName, TSDB_TABLE_FNAME_LEN, TOPIC_ENCODE_OVER);
362,571✔
128
  SDB_SET_INT32(pRaw, dataPos, pTopic->sqlLen, TOPIC_ENCODE_OVER);
362,571✔
129
  SDB_SET_BINARY(pRaw, dataPos, pTopic->sql, pTopic->sqlLen, TOPIC_ENCODE_OVER);
362,571✔
130
  SDB_SET_INT32(pRaw, dataPos, physicalPlanLen, TOPIC_ENCODE_OVER);
362,571✔
131
  if (physicalPlanLen) {
362,571✔
132
    SDB_SET_BINARY(pRaw, dataPos, pTopic->physicalPlan, physicalPlanLen, TOPIC_ENCODE_OVER);
362,571✔
133
  }
134
  SDB_SET_INT32(pRaw, dataPos, schemaLen, TOPIC_ENCODE_OVER);
362,571✔
135
  if (schemaLen) {
362,571✔
136
    swBuf = taosMemoryMalloc(schemaLen);
277,969✔
137
    if (swBuf == NULL) {
277,969✔
138
      goto TOPIC_ENCODE_OVER;
×
139
    }
140
    void *aswBuf = swBuf;
277,969✔
141
    if (taosEncodeSSchemaWrapper(&aswBuf, &pTopic->schema) < 0) {
555,938✔
142
      goto TOPIC_ENCODE_OVER;
×
143
    }
144
    SDB_SET_BINARY(pRaw, dataPos, swBuf, schemaLen, TOPIC_ENCODE_OVER);
277,969✔
145
  }
146
  SDB_SET_INT64(pRaw, dataPos, pTopic->ownerId, TOPIC_ENCODE_OVER); // since ver 4
362,571✔
147
  SDB_SET_RESERVE(pRaw, dataPos, MND_TOPIC_RESERVE_SIZE, TOPIC_ENCODE_OVER);
362,571✔
148
  SDB_SET_DATALEN(pRaw, dataPos, TOPIC_ENCODE_OVER);
362,571✔
149

150
  terrno = TSDB_CODE_SUCCESS;
362,571✔
151

152
TOPIC_ENCODE_OVER:
362,571✔
153
  if (swBuf) taosMemoryFree(swBuf);
362,571✔
154
  if (terrno != TSDB_CODE_SUCCESS) {
362,571✔
155
    mError("topic:%s, failed to encode to raw:%p since %s", pTopic->name, pRaw, terrstr());
×
156
    sdbFreeRaw(pRaw);
×
157
    return NULL;
×
158
  }
159

160
  mDebug("topic:%s, encode to raw:%p, row:%p", pTopic->name, pRaw, pTopic);
362,571✔
161
  return pRaw;
362,571✔
162
}
163

164
SSdbRow *mndTopicActionDecode(SSdbRaw *pRaw) {
308,723✔
165
  if (pRaw == NULL) return NULL;
308,723✔
166
  int32_t code = 0;
308,723✔
167
  int32_t lino = 0;
308,723✔
168
  terrno = TSDB_CODE_OUT_OF_MEMORY;
308,723✔
169
  SSdbRow *    pRow = NULL;
308,723✔
170
  SMqTopicObj *pTopic = NULL;
308,723✔
171
  void *       buf = NULL;
308,723✔
172
  char*        ast = NULL;
308,723✔
173

174
  int8_t sver = 0;
308,723✔
175
  if (sdbGetRawSoftVer(pRaw, &sver) != 0) goto TOPIC_DECODE_OVER;
308,723✔
176

177
  if (sver < 1 || sver > MND_TOPIC_VER_NUMBER) {
308,723✔
178
    terrno = TSDB_CODE_SDB_INVALID_DATA_VER;
×
179
    goto TOPIC_DECODE_OVER;
×
180
  }
181

182
  pRow = sdbAllocRow(sizeof(SMqTopicObj));
308,723✔
183
  if (pRow == NULL) goto TOPIC_DECODE_OVER;
308,723✔
184

185
  pTopic = sdbGetRowObj(pRow);
308,723✔
186
  if (pTopic == NULL) goto TOPIC_DECODE_OVER;
308,723✔
187

188
  int32_t len = 0;
308,723✔
189
  int32_t dataPos = 0;
308,723✔
190
  SDB_GET_BINARY(pRaw, dataPos, pTopic->name, TSDB_TOPIC_FNAME_LEN, TOPIC_DECODE_OVER);
308,723✔
191
  SDB_GET_BINARY(pRaw, dataPos, pTopic->db, TSDB_DB_FNAME_LEN, TOPIC_DECODE_OVER);
308,723✔
192
  if (sver >= 2) {
308,723✔
193
    SDB_GET_BINARY(pRaw, dataPos, pTopic->createUser, TSDB_USER_LEN, TOPIC_DECODE_OVER);
308,723✔
194
  }
195
  SDB_GET_INT64(pRaw, dataPos, &pTopic->createTime, TOPIC_DECODE_OVER);
308,723✔
196
  SDB_GET_INT64(pRaw, dataPos, &pTopic->updateTime, TOPIC_DECODE_OVER);
308,723✔
197
  SDB_GET_INT64(pRaw, dataPos, &pTopic->uid, TOPIC_DECODE_OVER);
308,723✔
198
  SDB_GET_INT64(pRaw, dataPos, &pTopic->dbUid, TOPIC_DECODE_OVER);
308,723✔
199
  SDB_GET_INT32(pRaw, dataPos, &pTopic->version, TOPIC_DECODE_OVER);
308,723✔
200
  SDB_GET_INT8(pRaw, dataPos, &pTopic->subType, TOPIC_DECODE_OVER);
308,723✔
201
  SDB_GET_INT8(pRaw, dataPos, &pTopic->withMeta, TOPIC_DECODE_OVER);
308,723✔
202

203
  SDB_GET_INT64(pRaw, dataPos, &pTopic->stbUid, TOPIC_DECODE_OVER);
308,723✔
204
  if (sver >= 3) {
308,723✔
205
    SDB_GET_BINARY(pRaw, dataPos, pTopic->stbName, TSDB_TABLE_FNAME_LEN, TOPIC_DECODE_OVER);
308,723✔
206
  }
207
  SDB_GET_INT32(pRaw, dataPos, &pTopic->sqlLen, TOPIC_DECODE_OVER);
308,723✔
208
  pTopic->sql = taosMemoryCalloc(pTopic->sqlLen, sizeof(char));
308,723✔
209
  if (pTopic->sql == NULL) {
308,723✔
210
    terrno = TSDB_CODE_OUT_OF_MEMORY;
×
211
    goto TOPIC_DECODE_OVER;
×
212
  }
213
  SDB_GET_BINARY(pRaw, dataPos, pTopic->sql, pTopic->sqlLen, TOPIC_DECODE_OVER);
308,723✔
214

215
  if (sver < MND_TOPIC_VER_SUPPORT_OWNER) {
308,723✔
216
    int32_t astLen = 0;
×
217
    SDB_GET_INT32(pRaw, dataPos, &astLen, TOPIC_DECODE_OVER);
×
218
    if (astLen) {
×
219
      ast = taosMemoryCalloc(astLen, sizeof(char));
×
220
      if (ast == NULL) {
×
221
        terrno = TSDB_CODE_OUT_OF_MEMORY;
×
222
        goto TOPIC_DECODE_OVER;
×
223
      }
224
      SDB_GET_BINARY(pRaw, dataPos, ast, astLen, TOPIC_DECODE_OVER);
×
225
      terrno = processAst(pTopic, ast);
×
226
      if (terrno != TSDB_CODE_SUCCESS) {
×
227
        goto TOPIC_DECODE_OVER;
×
228
      }
229
    }
230
  } else {
231
    SDB_GET_INT32(pRaw, dataPos, &len, TOPIC_DECODE_OVER);
308,723✔
232
    if (len) {
308,723✔
233
      pTopic->physicalPlan = taosMemoryCalloc(len, sizeof(char));
308,723✔
234
      if (pTopic->physicalPlan == NULL) {
308,723✔
235
        terrno = TSDB_CODE_OUT_OF_MEMORY;
×
236
        goto TOPIC_DECODE_OVER;
×
237
      }
238
      SDB_GET_BINARY(pRaw, dataPos, pTopic->physicalPlan, len, TOPIC_DECODE_OVER);
308,723✔
239
    } else {
240
      pTopic->physicalPlan = NULL;
×
241
    }
242

243
    SDB_GET_INT32(pRaw, dataPos, &len, TOPIC_DECODE_OVER);
308,723✔
244
    if (len) {
308,723✔
245
      buf = taosMemoryMalloc(len);
239,468✔
246
      if (buf == NULL) {
239,468✔
247
        terrno = TSDB_CODE_OUT_OF_MEMORY;
×
248
        goto TOPIC_DECODE_OVER;
×
249
      }
250
      SDB_GET_BINARY(pRaw, dataPos, buf, len, TOPIC_DECODE_OVER);
239,468✔
251
      if (taosDecodeSSchemaWrapper(buf, &pTopic->schema) == NULL) {
478,936✔
252
        goto TOPIC_DECODE_OVER;
×
253
      }
254
    } else {
255
      pTopic->schema.nCols = 0;
69,255✔
256
      pTopic->schema.version = 0;
69,255✔
257
      pTopic->schema.pSchema = NULL;
69,255✔
258
    }
259
    SDB_GET_INT64(pRaw, dataPos, &pTopic->ownerId, TOPIC_DECODE_OVER);
308,723✔
260
  }
261

262
  SDB_GET_RESERVE(pRaw, dataPos, MND_TOPIC_RESERVE_SIZE, TOPIC_DECODE_OVER);
308,723✔
263
  terrno = TSDB_CODE_SUCCESS;
308,723✔
264

265
TOPIC_DECODE_OVER:
308,723✔
266
  taosMemoryFreeClear(buf);
308,723✔
267
  taosMemoryFreeClear(ast);
308,723✔
268

269
  if (terrno != TSDB_CODE_SUCCESS) {
308,723✔
270
    mError("topic:%s, failed to decode from raw:%p since %s", pTopic == NULL ? "null" : pTopic->name, pRaw, terrstr());
×
271
    taosMemoryFreeClear(pRow);
×
272
    return NULL;
×
273
  }
274

275
  mDebug("topic:%s, decode from raw:%p, row:%p", pTopic->name, pRaw, pTopic);
308,723✔
276
  return pRow;
308,723✔
277
}
278

279
static int32_t mndTopicActionInsert(SSdb *pSdb, SMqTopicObj *pTopic) {
183,610✔
280
  mDebug("topic:%s perform insert action", pTopic != NULL ? pTopic->name : "null");
183,610✔
281
  return 0;
183,610✔
282
}
283

284
static int32_t mndTopicActionDelete(SSdb *pSdb, SMqTopicObj *pTopic) {
308,723✔
285
  if (pTopic == NULL) return 0;
308,723✔
286
  mDebug("%p topic:%s perform delete action", pTopic, pTopic->name);
308,723✔
287
  taosMemoryFreeClear(pTopic->sql);
308,723✔
288
  taosMemoryFreeClear(pTopic->physicalPlan);
308,723✔
289
  if (pTopic->schema.nCols) taosMemoryFreeClear(pTopic->schema.pSchema);
308,723✔
290
  return 0;
308,723✔
291
}
292

293
static int32_t mndTopicActionUpdate(SSdb *pSdb, SMqTopicObj *pOldTopic, SMqTopicObj *pNewTopic) {
2,492✔
294
  if (pOldTopic == NULL || pNewTopic == NULL) return 0;
2,492✔
295
  mDebug("topic:%s perform update action", pOldTopic->name);
2,492✔
296
  taosWLockLatch(&pOldTopic->lock);
2,492✔
297
  SMqTopicObj tmpTopic = *pOldTopic;
2,492✔
298
  (void)memcpy(pOldTopic, pNewTopic, offsetof(SMqTopicObj, lock));
2,492✔
299
  *pNewTopic = tmpTopic;
2,492✔
300
  taosWUnLockLatch(&pOldTopic->lock);
2,492✔
301
  return 0;
2,492✔
302
}
303

304
int32_t mndAcquireTopic(SMnode *pMnode, const char *topicName, SMqTopicObj **pTopic) {
2,842,079✔
305
  if (pMnode == NULL || topicName == NULL || pTopic == NULL) {
2,842,079✔
306
    return TSDB_CODE_INVALID_PARA;
×
307
  }
308
  SSdb *pSdb = pMnode->pSdb;
2,842,079✔
309
  *pTopic = sdbAcquire(pSdb, SDB_TOPIC, topicName);
2,842,079✔
310
  if (*pTopic == NULL) {
2,842,079✔
311
    return TSDB_CODE_MND_TOPIC_NOT_EXIST;
209,859✔
312
  }
313
  return TSDB_CODE_SUCCESS;
2,632,220✔
314
}
315

316
void mndReleaseTopic(SMnode *pMnode, SMqTopicObj *pTopic) {
3,141,277✔
317
  if (pMnode == NULL) return;
3,141,277✔
318
  SSdb *pSdb = pMnode->pSdb;
3,141,277✔
319
  sdbRelease(pSdb, pTopic);
3,141,277✔
320
}
321

322
static int32_t mndCheckCreateTopicReq(SCMCreateTopicReq *pCreate) {
163,531✔
323
  if (pCreate == NULL) return TSDB_CODE_INVALID_PARA;
163,531✔
324
  if (pCreate->sql == NULL) return TSDB_CODE_MND_INVALID_TOPIC;
163,531✔
325

326
  if (pCreate->subType == TOPIC_SUB_TYPE__COLUMN) {
163,531✔
327
    if (pCreate->ast == NULL || pCreate->ast[0] == 0) return TSDB_CODE_MND_INVALID_TOPIC;
118,338✔
328
  } else if (pCreate->subType == TOPIC_SUB_TYPE__TABLE) {
45,193✔
329
    if (pCreate->subStbName[0] == 0) return TSDB_CODE_MND_INVALID_TOPIC;
15,197✔
330
  } else if (pCreate->subType == TOPIC_SUB_TYPE__DB) {
29,996✔
331
    if (pCreate->subDbName[0] == 0) return TSDB_CODE_MND_INVALID_TOPIC;
29,996✔
332
  }
333

334
  return 0;
163,531✔
335
}
336

337
static int32_t processAst(SMqTopicObj *topicObj, const char *ast) {
159,723✔
338
  SNode *     pAst = NULL;
159,723✔
339
  SQueryPlan *pPlan = NULL;
159,723✔
340
  int32_t     code = TSDB_CODE_SUCCESS;
159,723✔
341
  int32_t     lino = 0;
159,723✔
342

343
  PRINT_LOG_START
159,723✔
344
  if (ast == NULL) {
159,723✔
345
    topicObj->physicalPlan = taosStrdup("");
36,446✔
346
    goto END;
36,446✔
347
  }
348
  qDebugL("%s topic:%s ast %s", __func__, topicObj->name, ast);
123,277✔
349
  MND_TMQ_RETURN_CHECK(nodesStringToNode(ast, &pAst));
123,277✔
350
  MND_TMQ_RETURN_CHECK(qExtractResultSchema(pAst, &topicObj->schema.nCols, &topicObj->schema.pSchema));
123,277✔
351

352
  SPlanContext cxt = {.pAstRoot = pAst, .topicQuery = true};
123,277✔
353
  MND_TMQ_RETURN_CHECK(qCreateQueryPlan(&cxt, &pPlan, NULL));
123,277✔
354
  if (pPlan == NULL) {
123,277✔
355
    code = TSDB_CODE_MND_INVALID_TOPIC_QUERY;
×
356
    goto END;
×
357
  }
358
  int32_t levelNum = LIST_LENGTH(pPlan->pSubplans);
123,277✔
359
  if (levelNum != 1) {
123,277✔
360
    code = TSDB_CODE_MND_INVALID_TOPIC_QUERY;
×
361
    goto END;
×
362
  }
363

364
  SNodeListNode *pNodeListNode = (SNodeListNode *)nodesListGetNode(pPlan->pSubplans, 0);
123,277✔
365
  MND_TMQ_NULL_CHECK(pNodeListNode);
123,277✔
366
  int32_t opNum = LIST_LENGTH(pNodeListNode->pNodeList);
123,277✔
367
  if (opNum != 1) {
123,277✔
368
    code = TSDB_CODE_MND_INVALID_TOPIC_QUERY;
×
369
    goto END;
×
370
  }
371

372
  code = nodesNodeToString(nodesListGetNode(pNodeListNode->pNodeList, 0), false, &topicObj->physicalPlan, NULL);
123,277✔
373

374
END:
159,723✔
375
  nodesDestroyNode(pAst);
159,723✔
376
  qDestroyQueryPlan(pPlan);
159,723✔
377
  PRINT_LOG_END
159,723✔
378
  return code;
159,723✔
379
}
380

381
static int32_t mndCreateTopic(SMnode *pMnode, SRpcMsg *pReq, SCMCreateTopicReq *pCreate, SDbObj *pDb,
157,231✔
382
                              SUserObj *pOperUser) {
383
  if (pMnode == NULL || pReq == NULL || pCreate == NULL || pDb == NULL || pOperUser == NULL)
157,231✔
384
    return TSDB_CODE_INVALID_PARA;
×
385
  STrans *    pTrans = NULL;
157,231✔
386
  int32_t     code = 0;
157,231✔
387
  int32_t     lino = 0;
157,231✔
388
  SMqTopicObj topicObj = {0};
157,231✔
389

390
  PRINT_LOG_START
157,231✔
391
  mInfo("start to create topic:%s", pCreate->name);
157,231✔
392
  pTrans = mndTransCreate(pMnode, TRN_POLICY_RETRY, TRN_CONFLICT_DB, pReq, "create-topic");
157,231✔
393
  MND_TMQ_NULL_CHECK(pTrans);
157,231✔
394
  mndTransSetDbName(pTrans, pDb->name, NULL);
157,231✔
395
  MND_TMQ_RETURN_CHECK(mndTransCheckConflict(pMnode, pTrans));
157,231✔
396

397
  tstrncpy(topicObj.name, pCreate->name, TSDB_TOPIC_FNAME_LEN);
157,231✔
398
  tstrncpy(topicObj.db, pDb->name, TSDB_DB_FNAME_LEN);
157,231✔
399
  tstrncpy(topicObj.createUser, pOperUser->name, TSDB_USER_LEN);
157,231✔
400
  topicObj.ownerId = pOperUser->uid;
157,231✔
401

402
  // MND_TMQ_RETURN_CHECK(mndCheckTopicPrivilege(pMnode, RPC_MSG_USER(pReq), MND_OPER_CREATE_TOPIC, &topicObj));
403
  if (pDb) {
157,231✔
404
    // already checked in parser, just check db use privilege here
405
    MND_TMQ_RETURN_CHECK(
157,231✔
406
        mndCheckDbPrivilege(pMnode, RPC_MSG_USER(pReq), RPC_MSG_TOKEN(pReq), MND_OPER_CREATE_TOPIC, pDb));
407
  }
408

409
  topicObj.createTime = taosGetTimestampMs();
157,231✔
410
  topicObj.updateTime = topicObj.createTime;
157,231✔
411
  topicObj.uid = mndGenerateUid(pCreate->name, strlen(pCreate->name));
157,231✔
412
  topicObj.dbUid = pDb->uid;
157,231✔
413
  topicObj.version = 1;
157,231✔
414
  topicObj.sql = taosStrdup(pCreate->sql);
157,231✔
415
  MND_TMQ_NULL_CHECK(topicObj.sql);
157,231✔
416
  topicObj.sqlLen = strlen(pCreate->sql) + 1;
157,231✔
417
  topicObj.subType = pCreate->subType;
157,231✔
418
  topicObj.withMeta = pCreate->withMeta;
157,231✔
419
  taosInitRWLatch(&topicObj.lock);
157,231✔
420

421
  MND_TMQ_RETURN_CHECK(processAst(&topicObj, pCreate->ast));
157,231✔
422

423
  if (pCreate->subStbName[0] != 0) {
157,231✔
424
    tstrncpy(topicObj.stbName, pCreate->subStbName, TSDB_TABLE_FNAME_LEN);
116,844✔
425
    SStbObj *pStb = mndAcquireStb(pMnode, topicObj.stbName);
116,844✔
426
    MND_TMQ_NULL_CHECK(pStb);
116,844✔
427
    char stbName[TSDB_TABLE_NAME_LEN] = {0};
116,076✔
428
    mndExtractTbNameFromStbFullName(pStb->name, stbName, TSDB_TABLE_NAME_LEN);
116,076✔
429
    MND_TMQ_RETURN_CHECK(
116,076✔
430
        mndCheckObjPrivilegeRecF(pMnode, pOperUser, PRIV_TOPIC_CREATE, PRIV_OBJ_DB, pStb->ownerId, pDb->name, NULL));
431
    MND_TMQ_RETURN_CHECK(
116,076✔
432
        mndCheckObjPrivilegeRecF(pMnode, pOperUser, PRIV_TBL_SELECT, PRIV_OBJ_TBL, pStb->ownerId, pDb->name, stbName));
433
    topicObj.stbUid = pStb->uid;
116,076✔
434
    mndReleaseStb(pMnode, pStb);
116,076✔
435
  }
436

437
  if (pCreate->subType == TOPIC_SUB_TYPE__DB) {
156,463✔
438
    MND_TMQ_RETURN_CHECK(
28,820✔
439
        mndCheckObjPrivilegeRecF(pMnode, pOperUser, PRIV_TOPIC_CREATE, PRIV_OBJ_DB, pDb->ownerId, pDb->name, NULL));
440
    MND_TMQ_RETURN_CHECK(mndCheckObjPrivilegeRecF(pMnode, pOperUser, PRIV_TBL_SELECT, PRIV_OBJ_TBL, 0, pDb->name, "*"));
28,820✔
441
  } else if (pCreate->subType == TOPIC_SUB_TYPE__COLUMN) {
127,643✔
442
    // TODO: check privilege on table
443
  }
444

445
  SSdbRaw *pCommitRaw = mndTopicActionEncode(&topicObj);
156,463✔
446
  MND_TMQ_NULL_CHECK(pCommitRaw);
156,463✔
447
  code = mndTransAppendCommitlog(pTrans, pCommitRaw);
156,463✔
448
  if (code != 0) {
156,463✔
449
    sdbFreeRaw(pCommitRaw);
×
450
    goto END;
×
451
  }
452

453
  MND_TMQ_RETURN_CHECK(sdbSetRawStatus(pCommitRaw, SDB_STATUS_READY));
156,463✔
454
  MND_TMQ_RETURN_CHECK(mndTransPrepare(pMnode, pTrans));
156,463✔
455

456
END:
157,231✔
457
  PRINT_LOG_END
157,231✔
458
  taosMemoryFreeClear(topicObj.sql);
157,231✔
459
  taosMemoryFreeClear(topicObj.physicalPlan);
157,231✔
460
  if (topicObj.schema.nCols) {
157,231✔
461
    taosMemoryFreeClear(topicObj.schema.pSchema);
120,785✔
462
  }
463
  mndTransDrop(pTrans);
157,231✔
464
  return code;
157,231✔
465
}
466

467
static int32_t mndReloadTopic(SMnode *pMnode, SRpcMsg *pReq, SCMCreateTopicReq *pCreate, SDbObj *pDb,
2,492✔
468
                              const char *userName, SMqTopicObj *topicObjOri) {
469
  if (pMnode == NULL || pReq == NULL || pCreate == NULL || pDb == NULL || userName == NULL)
2,492✔
470
    return TSDB_CODE_INVALID_PARA;
×
471
  STrans *    pTrans = NULL;
2,492✔
472
  int32_t     code = 0;
2,492✔
473
  int32_t     lino = 0;
2,492✔
474
  SMqTopicObj topicObj = {0};
2,492✔
475

476
  PRINT_LOG_START
2,492✔
477
  mInfo("start to reload topic:%s", pCreate->name);
2,492✔
478
  pTrans = mndTransCreate(pMnode, TRN_POLICY_RETRY, TRN_CONFLICT_DB, pReq, "reload-topic");
2,492✔
479
  MND_TMQ_NULL_CHECK(pTrans);
2,492✔
480
  mndTransSetDbName(pTrans, pDb->name, NULL);
2,492✔
481
  MND_TMQ_RETURN_CHECK(mndTransCheckConflict(pMnode, pTrans));
2,492✔
482

483
  tstrncpy(topicObj.name, pCreate->name, TSDB_TOPIC_FNAME_LEN);
2,492✔
484
  tstrncpy(topicObj.db, pDb->name, TSDB_DB_FNAME_LEN);
2,492✔
485
  tstrncpy(topicObj.createUser, userName, TSDB_USER_LEN);
2,492✔
486

487
  MND_TMQ_RETURN_CHECK(mndCheckTopicPrivilege(pMnode, RPC_MSG_USER(pReq), RPC_MSG_TOKEN(pReq), MND_OPER_CREATE_TOPIC, &topicObj));
2,492✔
488

489
  taosRLockLatch(&topicObjOri->lock);
2,492✔
490
  topicObj.createTime = topicObjOri->createTime;
2,492✔
491
  topicObj.updateTime = taosGetTimestampMs();
2,492✔
492
  topicObj.uid = topicObjOri->uid;
2,492✔
493
  topicObj.dbUid = pDb->uid;
2,492✔
494
  topicObj.version = topicObjOri->version + 1;
2,492✔
495
  topicObj.sql = taosStrdup(pCreate->sql);
2,492✔
496
  topicObj.sqlLen = strlen(pCreate->sql) + 1;
2,492✔
497
  topicObj.subType = pCreate->subType;
2,492✔
498
  topicObj.withMeta = pCreate->withMeta;
2,492✔
499
  taosInitRWLatch(&topicObj.lock);
2,492✔
500
  taosRUnLockLatch(&topicObjOri->lock);
2,492✔
501

502
  MND_TMQ_RETURN_CHECK(processAst(&topicObj, pCreate->ast));
2,492✔
503

504
  if (pCreate->subStbName[0] != 0) {
2,492✔
505
    tstrncpy(topicObj.stbName, pCreate->subStbName, TSDB_TABLE_FNAME_LEN);
1,424✔
506
    SStbObj *pStb = mndAcquireStb(pMnode, topicObj.stbName);
1,424✔
507
    MND_TMQ_NULL_CHECK(pStb);
1,424✔
508
    topicObj.stbUid = pStb->uid;
1,424✔
509
    mndReleaseStb(pMnode, pStb);
1,424✔
510
  }
511

512
  SSdbRaw *pCommitRaw = mndTopicActionEncode(&topicObj);
2,492✔
513
  MND_TMQ_NULL_CHECK(pCommitRaw);
2,492✔
514
  code = mndTransAppendCommitlog(pTrans, pCommitRaw);
2,492✔
515
  if (code != 0) {
2,492✔
516
    sdbFreeRaw(pCommitRaw);
×
517
    goto END;
×
518
  }
519

520
  MND_TMQ_RETURN_CHECK(sdbSetRawStatus(pCommitRaw, SDB_STATUS_READY));
2,492✔
521
  MND_TMQ_RETURN_CHECK(mndTransPrepare(pMnode, pTrans));
2,492✔
522

523
END:
2,492✔
524
  PRINT_LOG_END
2,492✔
525
  taosMemoryFreeClear(topicObj.sql);
2,492✔
526
  taosMemoryFreeClear(topicObj.physicalPlan);
2,492✔
527
  if (topicObj.schema.nCols) {
2,492✔
528
    taosMemoryFreeClear(topicObj.schema.pSchema);
2,492✔
529
  }
530
  mndTransDrop(pTrans);
2,492✔
531
  return code;
2,492✔
532
}
533

534
static int32_t creatTopic(SRpcMsg *pReq, SCMCreateTopicReq *createTopicReq, SUserObj *pOperUser) {
160,327✔
535
  SMqTopicObj *pTopic = NULL;
160,327✔
536
  SDbObj *     pDb = NULL;
160,327✔
537
  int32_t      code = TSDB_CODE_SUCCESS;
160,327✔
538
  int32_t      lino = 0;
160,327✔
539
  SMnode *     pMnode = pReq->info.node;
160,327✔
540
  int64_t      tss = taosGetTimestampMs();
160,327✔
541

542
  PRINT_LOG_START
160,327✔
543
  mInfo("topic:%s start to create, sql:%s", createTopicReq->name, createTopicReq->sql);
160,327✔
544
  code = mndAcquireTopic(pMnode, createTopicReq->name, &pTopic);
160,327✔
545
  if (code == TSDB_CODE_SUCCESS) {
160,327✔
546
    mndReleaseTopic(pMnode, pTopic);
696✔
547
    if (createTopicReq->igExists) {
696✔
548
      mInfo("topic:%s already exist, ignore exist is set", createTopicReq->name);
696✔
549
    } else {
550
      code = TSDB_CODE_MND_TOPIC_ALREADY_EXIST;
×
551
    }
552
    goto END;
696✔
553
  } else if (code != TSDB_CODE_MND_TOPIC_NOT_EXIST) {
159,631✔
554
    goto END;
×
555
  }
556

557
  pDb = mndAcquireDb(pMnode, createTopicReq->subDbName);
159,631✔
558
  MND_TMQ_NULL_CHECK(pDb);
159,631✔
559

560
  if (pDb->cfg.walRetentionPeriod == 0) {
158,863✔
561
    code = TSDB_CODE_MND_DB_RETENTION_PERIOD_ZERO;
×
562
    goto END;
×
563
  }
564

565
  if (sdbGetSize(pMnode->pSdb, SDB_TOPIC) >= tmqMaxTopicNum) {
158,863✔
566
    code = TSDB_CODE_TMQ_TOPIC_OUT_OF_RANGE;
1,632✔
567
    goto END;
1,632✔
568
  }
569

570
  MND_TMQ_RETURN_CHECK(grantCheck(TSDB_GRANT_SUBSCRIPTION));
157,231✔
571
  MND_TMQ_RETURN_CHECK(mndCreateTopic(pMnode, pReq, createTopicReq, pDb, pOperUser));
157,231✔
572
  if (tsAuditLevel >= AUDIT_LEVEL_DATABASE) {
156,463✔
573
    int64_t tse = taosGetTimestampMs();
156,463✔
574
    double  duration = (double)(tse - tss);
156,463✔
575
    duration = duration / 1000;
156,463✔
576
    auditRecord(pReq, pMnode->clusterId, "createTopic", createTopicReq->subDbName, createTopicReq->name,
156,463✔
577
                createTopicReq->sql, strlen(createTopicReq->sql), duration, 0);
156,463✔
578
  }
579
  code = TSDB_CODE_ACTION_IN_PROGRESS;
156,463✔
580

581
END:
160,327✔
582
  if (code != 0 && code != TSDB_CODE_ACTION_IN_PROGRESS) {
160,327✔
583
    mError("%s failed, topic:%s since %s", __func__, createTopicReq->name, tstrerror(code));
3,168✔
584
  } else {
585
    mInfo("topic:%s create successfully", createTopicReq->name);
157,159✔
586
  }
587
  mndReleaseDb(pMnode, pDb);
160,327✔
588
  return code;
160,327✔
589
}
590

591
static int32_t reloadTopic(SRpcMsg *pReq, SCMCreateTopicReq *createTopicReq) {
3,204✔
592
  SMnode *     pMnode = pReq->info.node;
3,204✔
593
  int32_t      code = TSDB_CODE_SUCCESS;
3,204✔
594
  int32_t      lino = 0;
3,204✔
595
  SDbObj *     pDb = NULL;
3,204✔
596
  SMqTopicObj *pTopic = NULL;
3,204✔
597
  int64_t      tss = taosGetTimestampMs();
3,204✔
598

599
  PRINT_LOG_START
3,204✔
600
  code = mndAcquireTopic(pMnode, createTopicReq->name, &pTopic);
3,204✔
601
  if (code != 0) {
3,204✔
602
    if (createTopicReq->igExists) {
712✔
603
      mInfo("topic:%s, not exist, ignore not exist is set", createTopicReq->name);
356✔
604
      code = 0;
356✔
605
      goto END;
356✔
606
    } else {
607
      mError("topic:%s, failed to reload since %s", createTopicReq->name, tstrerror(code));
356✔
608
      goto END;
356✔
609
    }
610
  }
611

612
  pDb = mndAcquireDb(pMnode, createTopicReq->subDbName);
2,492✔
613
  MND_TMQ_NULL_CHECK(pDb);
2,492✔
614

615
  MND_TMQ_RETURN_CHECK(grantCheck(TSDB_GRANT_SUBSCRIPTION));
2,492✔
616
  MND_TMQ_RETURN_CHECK(mndReloadTopic(pMnode, pReq, createTopicReq, pDb, RPC_MSG_USER(pReq), pTopic));
2,492✔
617

618
  if (tsAuditLevel >= AUDIT_LEVEL_DATABASE) {
2,492✔
619
    int64_t tse = taosGetTimestampMs();
2,492✔
620
    double  duration = (double)(tse - tss);
2,492✔
621
    duration = duration / 1000;
2,492✔
622
    auditRecord(pReq, pMnode->clusterId, "reloadTopic", createTopicReq->subDbName, createTopicReq->name,
2,492✔
623
                createTopicReq->sql, strlen(createTopicReq->sql), duration, 0);
2,492✔
624
  }
625

626
  code = TSDB_CODE_ACTION_IN_PROGRESS;
2,492✔
627

628
  if (topicsToReload == NULL) {
2,492✔
629
    topicsToReload = taosHashInit(4, taosGetDefaultHashFunction(TSDB_DATA_TYPE_BINARY), true, HASH_NO_LOCK);
356✔
630
    MND_TMQ_NULL_CHECK(topicsToReload);
356✔
631
  }
632
  MND_TMQ_RETURN_CHECK(
2,492✔
633
      taosHashPut(topicsToReload, createTopicReq->name, strlen(createTopicReq->name), createTopicReq->name, 1));
634
  mInfo("topic:%s, marked to reload", createTopicReq->name);
2,492✔
635

636
END:
3,204✔
637
  if (code != 0 && code != TSDB_CODE_ACTION_IN_PROGRESS) {
3,204✔
638
    mError("%s failed, topic:%s since %s", __func__, createTopicReq->name, tstrerror(code));
356✔
639
  } else {
640
    mInfo("topic:%s create successfully", createTopicReq->name);
2,848✔
641
  }
642
  mndReleaseTopic(pMnode, pTopic);
3,204✔
643
  mndReleaseDb(pMnode, pDb);
3,204✔
644

645
  return code;
3,204✔
646
}
647

648
static int32_t mndProcessCreateTopicReq(SRpcMsg *pReq) {
163,531✔
649
  if (pReq == NULL || pReq->contLen <= 0) {
163,531✔
650
    return TSDB_CODE_INVALID_MSG;
×
651
  }
652
  SMnode *pMnode = pReq->info.node;
163,531✔
653
  SUserObj *pOperUser = NULL;
163,531✔
654
  int32_t code = TSDB_CODE_SUCCESS;
163,531✔
655
  int32_t lino = 0;
163,531✔
656

657
  SCMCreateTopicReq createTopicReq = {0};
163,531✔
658

659
  PRINT_LOG_START
163,531✔
660
  MND_TMQ_RETURN_CHECK(tDeserializeSCMCreateTopicReq(pReq->pCont, pReq->contLen, &createTopicReq));
163,531✔
661

662
  if ((code = mndAcquireUser(pMnode, RPC_MSG_USER(pReq), &pOperUser)) != 0) goto END;
163,531✔
663

664
  mInfo("topic:%s start to create, sql:%s", createTopicReq.name, createTopicReq.sql);
163,531✔
665

666
  MND_TMQ_RETURN_CHECK(mndCheckCreateTopicReq(&createTopicReq));
163,531✔
667

668
  if (createTopicReq.reload) {
163,531✔
669
    MND_TMQ_RETURN_CHECK(reloadTopic(pReq, &createTopicReq));
3,204✔
670
  } else {
671
    MND_TMQ_RETURN_CHECK(creatTopic(pReq, &createTopicReq, pOperUser));
160,327✔
672
  }
673

674
END:
163,374✔
675
  tFreeSCMCreateTopicReq(&createTopicReq);
163,531✔
676
  mndReleaseUser(pMnode, pOperUser);
163,531✔
677
  return code;
163,531✔
678
}
679

680
static int32_t mndDropTopic(SMnode *pMnode, STrans *pTrans, SRpcMsg *pReq, SMqTopicObj *pTopic) {
121,633✔
681
  if (pMnode == NULL || pTrans == NULL || pReq == NULL || pTopic == NULL) {
121,633✔
682
    return TSDB_CODE_INVALID_MSG;
×
683
  }
684
  int32_t  code = 0;
121,633✔
685
  int32_t  lino = 0;
121,633✔
686
  SSdbRaw *pCommitRaw = NULL;
121,633✔
687
  PRINT_LOG_START
121,633✔
688
  char topicFName[TSDB_TOPIC_FNAME_LEN + 1] = {0};                       // 1.topic
121,633✔
689
  mndTopicGetShowName(pTopic->name, topicFName);
121,633✔
690
  char topicDbFName[TSDB_DB_NAME_LEN + TSDB_TOPIC_FNAME_LEN + 1] = {0};  // 1.db.topic
121,633✔
691
  (void)snprintf(topicDbFName, sizeof(topicDbFName), "%s.%s", pTopic->db, topicFName);
121,633✔
692
  MND_TMQ_RETURN_CHECK(mndUserRemoveTopic(pMnode, pTrans, topicDbFName));
121,633✔
693
  pCommitRaw = mndTopicActionEncode(pTopic);
121,633✔
694
  MND_TMQ_NULL_CHECK(pCommitRaw);
121,633✔
695
  code = mndTransAppendCommitlog(pTrans, pCommitRaw);
121,633✔
696
  if (code != 0) {
121,633✔
697
    sdbFreeRaw(pCommitRaw);
×
698
    goto END;
×
699
  }
700
  MND_TMQ_RETURN_CHECK(sdbSetRawStatus(pCommitRaw, SDB_STATUS_DROPPED));
121,633✔
701
  MND_TMQ_RETURN_CHECK(mndTransPrepare(pMnode, pTrans));
121,633✔
702

703
END:
121,633✔
704
  PRINT_LOG_END
121,633✔
705
  return code;
121,633✔
706
}
707

708
bool checkTopic(SArray *topics, char *topicName) {
117,168✔
709
  if (topics == NULL || topicName == NULL) {
117,168✔
710
    return false;
×
711
  }
712
  int32_t sz = taosArrayGetSize(topics);
117,168✔
713
  for (int32_t i = 0; i < sz; i++) {
117,926✔
714
    char *name = taosArrayGetP(topics, i);
1,949✔
715
    if (name && strcmp(name, topicName) == 0) {
1,949✔
716
      return true;
1,191✔
717
    }
718
  }
719
  return false;
115,977✔
720
}
721

722
static int32_t checkConsumer(STrans *pTrans, SMqConsumerObj *pConsumer, bool deleteConsumer, char *topicName) {
38,337✔
723
  int32_t         code = 0;
38,337✔
724
  int32_t         lino = 0;
38,337✔
725
  SMqConsumerObj *pConsumerNew = NULL;
38,337✔
726

727
  taosRLockLatch(&pConsumer->lock);
38,337✔
728
  bool found1 = checkTopic(pConsumer->assignedTopics, topicName);
38,337✔
729
  bool found2 = checkTopic(pConsumer->rebRemovedTopics, topicName);
38,337✔
730
  bool found3 = checkTopic(pConsumer->rebNewTopics, topicName);
38,337✔
731
  if (found1 || found2 || found3) {
38,337✔
732
    if (deleteConsumer) {
807✔
733
      MND_TMQ_RETURN_CHECK(tNewSMqConsumerObj(pConsumer->consumerId, pConsumer->cgroup, CONSUMER_CLEAR, NULL, NULL, &pConsumerNew));
807✔
734
      MND_TMQ_RETURN_CHECK(mndSetConsumerDropLogs(pTrans, pConsumerNew));
807✔
735
      tDeleteSMqConsumerObj(pConsumerNew);
807✔
736
      pConsumerNew = NULL;
807✔
737
    } else {
738
      mError("topic:%s, failed to drop since subscribed by consumer:0x%" PRIx64 ", in consumer group %s", topicName,
×
739
             pConsumer->consumerId, pConsumer->cgroup);
740
      code = TSDB_CODE_MND_TOPIC_SUBSCRIBED;
×
741
      goto END;
×
742
    }
743
  }
744
END:
38,337✔
745
  taosRUnLockLatch(&pConsumer->lock);
38,337✔
746
  tDeleteSMqConsumerObj(pConsumerNew);
38,337✔
747
  return code;
38,337✔
748
}
749

750
static int32_t mndCheckConsumerByTopic(SMnode *pMnode, STrans *pTrans, char *topicName, bool deleteConsumer) {
121,633✔
751
  if (pMnode == NULL || pTrans == NULL || topicName == NULL) {
121,633✔
752
    return TSDB_CODE_INVALID_MSG;
×
753
  }
754
  int32_t         code = 0;
121,633✔
755
  int32_t         lino = 0;
121,633✔
756
  SSdb *          pSdb = pMnode->pSdb;
121,633✔
757
  void *          pIter = NULL;
121,633✔
758
  SMqConsumerObj *pConsumer = NULL;
121,633✔
759

760
  PRINT_LOG_START
121,633✔
761
  while (1) {
762
    pIter = sdbFetch(pSdb, SDB_CONSUMER, pIter, (void **)&pConsumer);
159,970✔
763
    if (pIter == NULL) {
159,970✔
764
      break;
121,633✔
765
    }
766

767
    MND_TMQ_RETURN_CHECK(checkConsumer(pTrans, pConsumer, deleteConsumer, topicName));
38,337✔
768
    sdbRelease(pSdb, pConsumer);
38,337✔
769
  }
770

771
END:
121,633✔
772
  PRINT_LOG_END
121,633✔
773
  sdbRelease(pSdb, pConsumer);
121,633✔
774
  sdbCancelFetch(pSdb, pIter);
121,633✔
775
  return code;
121,633✔
776
}
777

778
static int32_t mndProcessDropTopicReq(SRpcMsg *pReq) {
131,385✔
779
  if (pReq == NULL) {
131,385✔
780
    return TSDB_CODE_INVALID_MSG;
×
781
  }
782
  SMnode *       pMnode = pReq->info.node;
131,385✔
783
  SMDropTopicReq dropReq = {0};
131,385✔
784
  int32_t        code = 0;
131,385✔
785
  int32_t        lino = 0;
131,385✔
786
  SMqTopicObj *  pTopic = NULL;
131,385✔
787
  STrans *       pTrans = NULL;
131,385✔
788
  SUserObj      *pOperUser = NULL;
131,385✔
789
  int64_t        tss = taosGetTimestampMs();
131,385✔
790

791
  PRINT_LOG_START
131,385✔
792
  MND_TMQ_RETURN_CHECK(tDeserializeSMDropTopicReq(pReq->pCont, pReq->contLen, &dropReq));
131,385✔
793

794
  code = mndAcquireTopic(pMnode, dropReq.name, &pTopic);
131,385✔
795
  if (code != 0) {
131,385✔
796
    if (dropReq.igNotExists) {
9,203✔
797
      mInfo("topic:%s, not exist, ignore not exist is set", dropReq.name);
9,203✔
798
      code = 0;
9,203✔
799
    }
800
    goto END;
9,203✔
801
  }
802
  taosRLockLatch(&pTopic->lock);
122,182✔
803

804
  pTrans = mndTransCreate(pMnode, TRN_POLICY_RETRY, TRN_CONFLICT_DB, pReq, "drop-topic");
122,182✔
805
  MND_TMQ_NULL_CHECK(pTrans);
122,182✔
806

807
  mndTransSetDbName(pTrans, pTopic->db, NULL);
122,182✔
808
  MND_TMQ_RETURN_CHECK(mndTransCheckConflict(pMnode, pTrans));
122,182✔
809
  mInfo("trans:%d, used to drop topic:%s, force:%d", pTrans->id, pTopic->name, dropReq.force);
122,182✔
810

811
  MND_TMQ_RETURN_CHECK(mndAcquireUser(pMnode, RPC_MSG_USER(pReq), &pOperUser));
122,182✔
812

813
  // MND_TMQ_RETURN_CHECK(mndCheckTopicPrivilege(pMnode, RPC_MSG_USER(pReq), MND_OPER_DROP_TOPIC, pTopic));
814
  // MND_TMQ_RETURN_CHECK(mndCheckDbPrivilegeByName(pMnode, RPC_MSG_USER(pReq), MND_OPER_READ_DB, pTopic->db));
815
  MND_TMQ_RETURN_CHECK(
122,182✔
816
      mndCheckDbPrivilegeByName(pMnode, RPC_MSG_USER(pReq), RPC_MSG_TOKEN(pReq), MND_OPER_USE_DB, pTopic->db, true));
817
  MND_TMQ_RETURN_CHECK(mndCheckObjPrivilegeRecF(pMnode, pOperUser, PRIV_CM_DROP, PRIV_OBJ_TOPIC, pTopic->ownerId,
122,010✔
818
                                                pTopic->db, mndGetDbStr(pTopic->name)));
819
  MND_TMQ_RETURN_CHECK(mndCheckConsumerByTopic(pMnode, pTrans, dropReq.name, dropReq.force));
121,633✔
820
  MND_TMQ_RETURN_CHECK(mndDropSubByTopic(pMnode, pTrans, dropReq.name, dropReq.force));
121,633✔
821
  MND_TMQ_RETURN_CHECK(mndDropTopic(pMnode, pTrans, pReq, pTopic));
121,633✔
822
  if (tsAuditLevel >= AUDIT_LEVEL_DATABASE) {
121,633✔
823
    int64_t tse = taosGetTimestampMs();
121,633✔
824
    double  duration = (double)(tse - tss);
121,633✔
825
    duration = duration / 1000;
121,633✔
826
    auditRecord(pReq, pMnode->clusterId, "dropTopic", pTopic->db, dropReq.name, dropReq.sql, dropReq.sqlLen, duration,
121,633✔
827
                0);
828
  }
829

830
  code = TSDB_CODE_ACTION_IN_PROGRESS;
121,633✔
831

832
END:
131,385✔
833
  if (code != 0 && code != TSDB_CODE_ACTION_IN_PROGRESS) {
131,385✔
834
    mError("%s failed, topic:%s since %s", __func__, dropReq.name, tstrerror(code));
549✔
835
  } else {
836
    mInfo("topic:%s dropped successfully", dropReq.name);
130,836✔
837
  }
838
  if (pTopic != NULL) {
131,385✔
839
    taosRUnLockLatch(&pTopic->lock);
122,182✔
840
  }
841
  mndReleaseTopic(pMnode, pTopic);
131,385✔
842
  mndReleaseUser(pMnode, pOperUser);
131,385✔
843
  mndTransDrop(pTrans);
131,385✔
844
  tFreeSMDropTopicReq(&dropReq);
131,385✔
845
  return code;
131,385✔
846
}
847

848
int32_t mndGetNumOfTopics(SMnode *pMnode, char *dbName, int32_t *pNumOfTopics) {
130,901✔
849
  if (pMnode == NULL || dbName == NULL || pNumOfTopics == NULL) {
130,901✔
850
    return TSDB_CODE_INVALID_MSG;
×
851
  }
852
  *pNumOfTopics = 0;
130,901✔
853

854
  SSdb *  pSdb = pMnode->pSdb;
130,901✔
855
  SDbObj *pDb = mndAcquireDb(pMnode, dbName);
130,901✔
856
  if (pDb == NULL) {
130,901✔
857
    return TSDB_CODE_MND_DB_NOT_SELECTED;
×
858
  }
859

860
  int32_t numOfTopics = 0;
130,901✔
861
  void *  pIter = NULL;
130,901✔
862
  while (1) {
×
863
    SMqTopicObj *pTopic = NULL;
130,901✔
864
    pIter = sdbFetch(pSdb, SDB_TOPIC, pIter, (void **)&pTopic);
130,901✔
865
    if (pIter == NULL) {
130,901✔
866
      break;
130,901✔
867
    }
868
    taosRLockLatch(&pTopic->lock);
×
869
    if (pTopic->dbUid == pDb->uid) {
×
870
      numOfTopics++;
×
871
    }
872
    taosRUnLockLatch(&pTopic->lock);
×
873

874
    sdbRelease(pSdb, pTopic);
×
875
  }
876

877
  *pNumOfTopics = numOfTopics;
130,901✔
878
  mndReleaseDb(pMnode, pDb);
130,901✔
879
  return 0;
130,901✔
880
}
881

882
static void schemaToJson(SSchema *schema, int32_t nCols, char *schemaJson) {
57,416✔
883
  if (schema == NULL || schemaJson == NULL) {
57,416✔
884
    return;
×
885
  }
886
  char *  string = NULL;
57,416✔
887
  int32_t code = 0;
57,416✔
888
  int32_t lino = 0;
57,416✔
889

890
  cJSON *cbytes = NULL;
57,416✔
891
  cJSON *ctype = NULL;
57,416✔
892
  cJSON *cname = NULL;
57,416✔
893
  cJSON *column = NULL;
57,416✔
894
  cJSON *columns = cJSON_CreateArray();
57,416✔
895
  MND_TMQ_NULL_CHECK(columns);
57,416✔
896
  for (int i = 0; i < nCols; i++) {
643,343✔
897
    column = cJSON_CreateObject();
585,927✔
898
    MND_TMQ_NULL_CHECK(column);
585,927✔
899
    SSchema *s = schema + i;
585,927✔
900
    cname = cJSON_CreateString(s->name);
585,927✔
901
    MND_TMQ_NULL_CHECK(cname);
585,927✔
902
    MND_TMQ_CONDITION_CHECK(cJSON_AddItemToObject(column, "name", cname), 0);
585,927✔
903
    cname = NULL;  // ownership transferred to column object
585,927✔
904

905
    ctype = cJSON_CreateString(tDataTypes[s->type].name);
585,927✔
906
    MND_TMQ_NULL_CHECK(ctype);
585,927✔
907
    MND_TMQ_CONDITION_CHECK(cJSON_AddItemToObject(column, "type", ctype), 0);
585,927✔
908
    ctype = NULL;  // ownership transferred to column object
585,927✔
909

910
    int32_t length = 0;
585,927✔
911
    if (s->type == TSDB_DATA_TYPE_BINARY || s->type == TSDB_DATA_TYPE_VARBINARY || s->type == TSDB_DATA_TYPE_GEOMETRY) {
585,927✔
912
      length = s->bytes - VARSTR_HEADER_SIZE;
110,114✔
913
    } else if (s->type == TSDB_DATA_TYPE_NCHAR || s->type == TSDB_DATA_TYPE_JSON) {
475,813✔
914
      length = (s->bytes - VARSTR_HEADER_SIZE) / TSDB_NCHAR_SIZE;
73,054✔
915
    } else {
916
      length = s->bytes;
402,759✔
917
    }
918
    cbytes = cJSON_CreateNumber(length);
585,927✔
919
    MND_TMQ_NULL_CHECK(cbytes);
585,927✔
920
    MND_TMQ_CONDITION_CHECK(cJSON_AddItemToObject(column, "length", cbytes), 0);
585,927✔
921
    cbytes = NULL;  // ownership transferred to column object
585,927✔
922

923
    MND_TMQ_CONDITION_CHECK(cJSON_AddItemToArray(columns, column), 0);
585,927✔
924
    column = NULL;  // ownership transferred to columns array
585,927✔
925
  }
926

927
END:
57,416✔
928
  string = cJSON_PrintUnformatted(columns);
57,416✔
929
  cJSON_Delete(columns);
57,416✔
930
  cJSON_Delete(column);
57,416✔
931
  cJSON_Delete(cname);
57,416✔
932
  cJSON_Delete(ctype);
57,416✔
933
  cJSON_Delete(cbytes);
57,416✔
934

935
  size_t len = strlen(string);
57,416✔
936
  if (string && len <= TSDB_SHOW_SCHEMA_JSON_LEN) {
57,416✔
937
    STR_TO_VARSTR(schemaJson, string);
57,416✔
938
  } else {
939
    mError("mndRetrieveTopic build schema error json:%p, json len:%zu", string, len);
×
940
    STR_TO_VARSTR(schemaJson, "NULL");
×
941
  }
942
  taosMemoryFree(string);
57,416✔
943
}
944

945
static int32_t buildResult(SMqTopicObj *pTopic, int32_t *numOfRows, SMnode *pMnode, SSDataBlock *pBlock) {
59,092✔
946
  SColumnInfoData *pColInfo = NULL;
59,092✔
947
  SName            n = {0};
59,092✔
948
  int32_t          cols = 0;
59,092✔
949
  char *           schemaJson = NULL;
59,092✔
950
  char *           sql = NULL;
59,092✔
951
  int32_t          code = 0;
59,092✔
952
  int32_t          lino = 0;
59,092✔
953

954
  taosRLockLatch(&pTopic->lock);
59,092✔
955

956
  char        topicName[TSDB_TOPIC_NAME_LEN + VARSTR_HEADER_SIZE + 5] = {0};
59,092✔
957
  const char *pName = mndGetDbStr(pTopic->name);
59,092✔
958
  STR_TO_VARSTR(topicName, pName);
59,092✔
959

960
  pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
59,092✔
961
  MND_TMQ_NULL_CHECK(pColInfo);
59,092✔
962
  MND_TMQ_RETURN_CHECK(colDataSetVal(pColInfo, *numOfRows, (const char *)topicName, false));
59,092✔
963

964
  char dbName[TSDB_DB_NAME_LEN + VARSTR_HEADER_SIZE] = {0};
59,092✔
965
  MND_TMQ_RETURN_CHECK(tNameFromString(&n, pTopic->db, T_NAME_ACCT | T_NAME_DB));
59,092✔
966
  MND_TMQ_RETURN_CHECK(tNameGetDbName(&n, varDataVal(dbName)));
59,092✔
967
  varDataSetLen(dbName, strlen(varDataVal(dbName)));
59,092✔
968

969
  pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
59,092✔
970
  MND_TMQ_NULL_CHECK(pColInfo);
59,092✔
971
  MND_TMQ_RETURN_CHECK(colDataSetVal(pColInfo, *numOfRows, (const char *)dbName, false));
59,092✔
972

973
  pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
59,092✔
974
  MND_TMQ_NULL_CHECK(pColInfo);
59,092✔
975
  MND_TMQ_RETURN_CHECK(colDataSetVal(pColInfo, *numOfRows, (const char *)&pTopic->createTime, false));
59,092✔
976

977
  sql = taosMemoryMalloc(strlen(pTopic->sql) + VARSTR_HEADER_SIZE);
59,092✔
978
  MND_TMQ_NULL_CHECK(sql);
59,092✔
979
  STR_TO_VARSTR(sql, pTopic->sql);
59,092✔
980

981
  pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
59,092✔
982
  MND_TMQ_NULL_CHECK(pColInfo);
59,092✔
983
  MND_TMQ_RETURN_CHECK(colDataSetVal(pColInfo, *numOfRows, (const char *)sql, false));
59,092✔
984

985
  taosMemoryFreeClear(sql);
59,092✔
986

987
  schemaJson = taosMemoryMalloc(TSDB_SHOW_SCHEMA_JSON_LEN + VARSTR_HEADER_SIZE);
59,092✔
988
  MND_TMQ_NULL_CHECK(schemaJson);
59,092✔
989
  if (pTopic->subType == TOPIC_SUB_TYPE__COLUMN) {
59,092✔
990
    schemaToJson(pTopic->schema.pSchema, pTopic->schema.nCols, schemaJson);
57,009✔
991
  } else if (pTopic->subType == TOPIC_SUB_TYPE__TABLE) {
2,083✔
992
    SStbObj *pStb = mndAcquireStb(pMnode, pTopic->stbName);
407✔
993
    if (pStb == NULL) {
407✔
994
      STR_TO_VARSTR(schemaJson, "NULL");
×
995
      mError("mndRetrieveTopic mndAcquireStb null stbName:%s", pTopic->stbName);
×
996
    } else {
997
      schemaToJson(pStb->pColumns, pStb->numOfColumns, schemaJson);
407✔
998
      mndReleaseStb(pMnode, pStb);
407✔
999
    }
1000
  } else {
1001
    STR_TO_VARSTR(schemaJson, "NULL");
1,676✔
1002
  }
1003

1004
  pColInfo = taosArrayGet(pBlock->pDataBlock, cols++);
59,092✔
1005
  MND_TMQ_NULL_CHECK(pColInfo);
59,092✔
1006
  MND_TMQ_RETURN_CHECK(colDataSetVal(pColInfo, *numOfRows, (const char *)schemaJson, false));
59,092✔
1007
  taosMemoryFreeClear(schemaJson);
59,092✔
1008

1009
  (*numOfRows)++;
59,092✔
1010

1011
END:
59,092✔
1012
  taosRUnLockLatch(&pTopic->lock);
59,092✔
1013

1014
  taosMemoryFreeClear(sql);
59,092✔
1015
  taosMemoryFreeClear(schemaJson);
59,092✔
1016
  return code;
59,092✔
1017
}
1018

1019
static int32_t mndRetrieveTopic(SRpcMsg *pReq, SShowObj *pShow, SSDataBlock *pBlock, int32_t rowsCapacity) {
22,294✔
1020
  if (pReq == NULL || pShow == NULL || pBlock == NULL) {
22,294✔
1021
    return TSDB_CODE_INVALID_MSG;
×
1022
  }
1023
  SMnode      *pMnode = pReq->info.node;
22,294✔
1024
  SSdb        *pSdb = pMnode->pSdb;
22,294✔
1025
  int32_t      numOfRows = 0;
22,294✔
1026
  SMqTopicObj *pTopic = NULL;
22,294✔
1027
  SUserObj    *pOperUser = NULL;
22,294✔
1028
  int32_t      code = 0, lino = 0;
22,294✔
1029
  char        *sql = NULL;
22,294✔
1030
  char        *schemaJson = NULL;
22,294✔
1031
  char         objFName[TSDB_OBJ_FNAME_LEN + 1] = {0};
22,294✔
1032
  bool         showAll = false;
22,294✔
1033

1034
  MND_TMQ_RETURN_CHECK(mndAcquireUser(pMnode, RPC_MSG_USER(pReq), &pOperUser));
22,294✔
1035
  (void)snprintf(objFName, sizeof(objFName), "%d.*", pOperUser->acctId);
22,294✔
1036
  int32_t objLevel = privObjGetLevel(PRIV_OBJ_TOPIC);
22,294✔
1037
  showAll =
22,294✔
1038
      (0 == mndCheckSysObjPrivilege(pMnode, pOperUser, RPC_MSG_TOKEN(pReq), PRIV_CM_SHOW, PRIV_OBJ_TOPIC, 0, objFName,
22,294✔
1039
                                    objLevel == 0 ? NULL : "*"));  // 1.*.*
1040

1041
  PRINT_LOG_START
22,294✔
1042

1043
  while (numOfRows < rowsCapacity) {
81,902✔
1044
    pShow->pIter = sdbFetch(pSdb, SDB_TOPIC, pShow->pIter, (void **)&pTopic);
81,902✔
1045
    if (pShow->pIter == NULL) break;
81,902✔
1046

1047
    if (!showAll) {
59,608✔
1048
      if (mndCheckObjPrivilegeRecF(pMnode, pOperUser, PRIV_CM_SHOW, PRIV_OBJ_TOPIC, pTopic->ownerId, pTopic->db,
1,376✔
1049
                                   objLevel == 0 ? NULL : mndGetDbStr(pTopic->name)) != 0) {  // 1.topic1
688✔
1050
        sdbRelease(pSdb, pTopic);
516✔
1051
        continue;
516✔
1052
      }
1053
    }
1054

1055
    MND_TMQ_RETURN_CHECK(buildResult(pTopic, &numOfRows, pMnode, pBlock));
59,092✔
1056
    sdbRelease(pSdb, pTopic);
59,092✔
1057
    pTopic = NULL;
59,092✔
1058
  }
1059
  pShow->numOfRows += numOfRows;
22,294✔
1060

1061
END:
22,294✔
1062
  sdbCancelFetch(pSdb, pShow->pIter);
22,294✔
1063
  sdbRelease(pSdb, pTopic);
22,294✔
1064
  mndReleaseUser(pMnode, pOperUser);
22,294✔
1065
  if (code != TSDB_CODE_SUCCESS) {
22,294✔
1066
    mError("%s failed since %s", __func__, tstrerror(code));
×
1067
    return code;
×
1068
  } else {
1069
    mDebug("%s retrieved %d rows successfully", __func__, numOfRows);
22,294✔
1070
    return numOfRows;
22,294✔
1071
  }
1072
}
1073

1074
static void mndCancelGetNextTopic(SMnode *pMnode, void *pIter) {
×
1075
  if (pMnode == NULL) return;
×
1076
  SSdb *pSdb = pMnode->pSdb;
×
1077
  sdbCancelFetchByType(pSdb, pIter, SDB_TOPIC);
×
1078
}
1079

1080
bool mndTopicExistsForDb(SMnode *pMnode, SDbObj *pDb) {
823,110✔
1081
  if (pMnode == NULL || pDb == NULL) {
823,110✔
1082
    return false;
×
1083
  }
1084
  SSdb *       pSdb = pMnode->pSdb;
823,110✔
1085
  void *       pIter = NULL;
823,110✔
1086
  SMqTopicObj *pTopic = NULL;
823,110✔
1087

1088
  while (1) {
63,536✔
1089
    pIter = sdbFetch(pSdb, SDB_TOPIC, pIter, (void **)&pTopic);
886,646✔
1090
    if (pIter == NULL) {
886,646✔
1091
      break;
822,702✔
1092
    }
1093

1094
    taosRLockLatch(&pTopic->lock);
63,944✔
1095
    bool found = pTopic->dbUid == pDb->uid;
63,944✔
1096
    taosRUnLockLatch(&pTopic->lock);
63,944✔
1097

1098
    if (found) {
63,944✔
1099
      sdbRelease(pSdb, pTopic);
408✔
1100
      sdbCancelFetch(pSdb, pIter);
408✔
1101
      return true;
408✔
1102
    }
1103

1104
    sdbRelease(pSdb, pTopic);
63,536✔
1105
  }
1106

1107
  return false;
822,702✔
1108
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc