Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
taosdata
TDengine
提交
43083fa3
TDengine
项目概览
taosdata
/
TDengine
接近 2 年 前同步成功
通知
1192
Star
22018
Fork
4786
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
1
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
TDengine
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
1
Issue
1
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
提交
43083fa3
编写于
3月 17, 2022
作者:
S
Shengliang Guan
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
shm
上级
b16e6868
变更
6
隐藏空白更改
内联
并排
Showing
6 changed file
with
98 addition
and
61 deletion
+98
-61
source/dnode/mgmt/container/inc/dnd.h
source/dnode/mgmt/container/inc/dnd.h
+3
-0
source/dnode/mgmt/container/inc/dndInt.h
source/dnode/mgmt/container/inc/dndInt.h
+0
-1
source/dnode/mgmt/container/src/dndMsg.c
source/dnode/mgmt/container/src/dndMsg.c
+72
-0
source/dnode/mgmt/dnode/src/dmWorker.c
source/dnode/mgmt/dnode/src/dmWorker.c
+10
-46
source/dnode/mgmt/test/sut/inc/server.h
source/dnode/mgmt/test/sut/inc/server.h
+1
-1
source/dnode/mgmt/test/sut/src/server.cpp
source/dnode/mgmt/test/sut/src/server.cpp
+12
-13
未找到文件。
source/dnode/mgmt/container/inc/dnd.h
浏览文件 @
43083fa3
...
@@ -158,6 +158,9 @@ int32_t dndInitWorker(void *param, SDnodeWorker *pWorker, EWorkerType type, cons
...
@@ -158,6 +158,9 @@ int32_t dndInitWorker(void *param, SDnodeWorker *pWorker, EWorkerType type, cons
void
dndCleanupWorker
(
SDnodeWorker
*
pWorker
);
void
dndCleanupWorker
(
SDnodeWorker
*
pWorker
);
int32_t
dndWriteMsgToWorker
(
SDnodeWorker
*
pWorker
,
void
*
pCont
,
int32_t
contLen
);
int32_t
dndWriteMsgToWorker
(
SDnodeWorker
*
pWorker
,
void
*
pCont
,
int32_t
contLen
);
int32_t
dndProcessCreateNodeMsg
(
SDnode
*
pDnode
,
SNodeMsg
*
pMsg
);
int32_t
dndProcessDropNodeMsg
(
SDnode
*
pDnode
,
SNodeMsg
*
pMsg
);
#ifdef __cplusplus
#ifdef __cplusplus
}
}
#endif
#endif
...
...
source/dnode/mgmt/container/inc/dndInt.h
浏览文件 @
43083fa3
...
@@ -20,7 +20,6 @@
...
@@ -20,7 +20,6 @@
#include "bm.h"
#include "bm.h"
#include "dm.h"
#include "dm.h"
#include "dndInt.h"
#include "mm.h"
#include "mm.h"
#include "qmInt.h"
#include "qmInt.h"
#include "smInt.h"
#include "smInt.h"
...
...
source/dnode/mgmt/container/src/dndMsg.c
0 → 100644
浏览文件 @
43083fa3
/*
* 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
#include "dndInt.h"
static
SMgmtWrapper
*
dndGetWrapperFromMsg
(
SDnode
*
pDnode
,
SNodeMsg
*
pMsg
)
{
SMgmtWrapper
*
pWrapper
=
NULL
;
switch
(
pMsg
->
rpcMsg
.
msgType
)
{
case
TDMT_DND_CREATE_MNODE
:
return
dndGetWrapper
(
pDnode
,
MNODE
);
case
TDMT_DND_CREATE_QNODE
:
return
dndGetWrapper
(
pDnode
,
QNODE
);
case
TDMT_DND_CREATE_SNODE
:
return
dndGetWrapper
(
pDnode
,
SNODE
);
case
TDMT_DND_CREATE_BNODE
:
return
dndGetWrapper
(
pDnode
,
BNODE
);
default:
return
NULL
;
}
}
int32_t
dndProcessCreateNodeMsg
(
SDnode
*
pDnode
,
SNodeMsg
*
pMsg
)
{
SMgmtWrapper
*
pWrapper
=
dndGetWrapperFromMsg
(
pDnode
,
pMsg
);
if
(
pWrapper
->
procType
==
PROC_SINGLE
)
{
switch
(
pMsg
->
rpcMsg
.
msgType
)
{
case
TDMT_DND_CREATE_MNODE
:
return
mmProcessCreateReq
(
pWrapper
->
pMgmt
,
pMsg
);
case
TDMT_DND_CREATE_QNODE
:
return
qmProcessCreateReq
(
pWrapper
->
pMgmt
,
pMsg
);
case
TDMT_DND_CREATE_SNODE
:
return
smProcessCreateReq
(
pWrapper
->
pMgmt
,
pMsg
);
case
TDMT_DND_CREATE_BNODE
:
return
bmProcessCreateReq
(
pWrapper
->
pMgmt
,
pMsg
);
default:
terrno
=
TSDB_CODE_MSG_NOT_PROCESSED
;
return
-
1
;
}
}
else
{
terrno
=
TSDB_CODE_MSG_NOT_PROCESSED
;
return
-
1
;
}
}
int32_t
dndProcessDropNodeMsg
(
SDnode
*
pDnode
,
SNodeMsg
*
pMsg
)
{
SMgmtWrapper
*
pWrapper
=
dndGetWrapperFromMsg
(
pDnode
,
pMsg
);
switch
(
pMsg
->
rpcMsg
.
msgType
)
{
case
TDMT_DND_DROP_MNODE
:
return
mmProcessDropReq
(
pWrapper
->
pMgmt
,
pMsg
);
case
TDMT_DND_DROP_QNODE
:
return
qmProcessDropReq
(
pWrapper
->
pMgmt
,
pMsg
);
case
TDMT_DND_DROP_SNODE
:
return
smProcessDropReq
(
pWrapper
->
pMgmt
,
pMsg
);
case
TDMT_DND_DROP_BNODE
:
return
bmProcessDropReq
(
pWrapper
->
pMgmt
,
pMsg
);
default:
terrno
=
TSDB_CODE_MSG_NOT_PROCESSED
;
return
-
1
;
}
}
\ No newline at end of file
source/dnode/mgmt/dnode/src/dmWorker.c
浏览文件 @
43083fa3
...
@@ -53,62 +53,27 @@ static void *dmThreadRoutine(void *param) {
...
@@ -53,62 +53,27 @@ static void *dmThreadRoutine(void *param) {
}
}
static
void
dmProcessQueue
(
SDnodeMgmt
*
pMgmt
,
SNodeMsg
*
pMsg
)
{
static
void
dmProcessQueue
(
SDnodeMgmt
*
pMgmt
,
SNodeMsg
*
pMsg
)
{
int32_t
code
=
-
1
;
int32_t
code
=
-
1
;
tmsg_t
msgType
=
pMsg
->
rpcMsg
.
msgType
;
tmsg_t
msgType
=
pMsg
->
rpcMsg
.
msgType
;
SDnode
*
pDnode
=
pMgmt
->
pDnode
;
SDnode
*
pDnode
=
pMgmt
->
pDnode
;
SMgmtWrapper
*
pWrapper
=
NULL
;
dTrace
(
"msg:%p, will be processed"
,
pMsg
);
dTrace
(
"msg:%p, will be processed"
,
pMsg
);
switch
(
msgType
)
{
switch
(
msgType
)
{
case
TDMT_DND_CREATE_MNODE
:
case
TDMT_DND_CREATE_MNODE
:
pWrapper
=
dndGetWrapper
(
pDnode
,
MNODE
);
code
=
mmProcessCreateReq
(
pWrapper
->
pMgmt
,
pMsg
);
break
;
case
TDMT_DND_DROP_MNODE
:
pWrapper
=
dndGetWrapper
(
pDnode
,
MNODE
);
code
=
mmProcessDropReq
(
pWrapper
->
pMgmt
,
pMsg
);
break
;
case
TDMT_DND_CREATE_QNODE
:
case
TDMT_DND_CREATE_QNODE
:
pWrapper
=
dndGetWrapper
(
pDnode
,
QNODE
);
code
=
qmProcessCreateReq
(
pWrapper
->
pMgmt
,
pMsg
);
break
;
case
TDMT_DND_DROP_QNODE
:
pWrapper
=
dndGetWrapper
(
pDnode
,
QNODE
);
code
=
qmProcessDropReq
(
pWrapper
->
pMgmt
,
pMsg
);
break
;
case
TDMT_DND_CREATE_SNODE
:
case
TDMT_DND_CREATE_SNODE
:
pWrapper
=
dndGetWrapper
(
pDnode
,
SNODE
);
code
=
smProcessCreateReq
(
pWrapper
->
pMgmt
,
pMsg
);
break
;
case
TDMT_DND_DROP_SNODE
:
pWrapper
=
dndGetWrapper
(
pDnode
,
SNODE
);
code
=
smProcessDropReq
(
pWrapper
->
pMgmt
,
pMsg
);
break
;
case
TDMT_DND_CREATE_BNODE
:
case
TDMT_DND_CREATE_BNODE
:
pWrapper
=
dndGetWrapper
(
pDnode
,
BNODE
);
code
=
dndProcessCreateNodeMsg
(
pMgmt
->
pDnode
,
pMsg
);
code
=
bmProcessCreateReq
(
pWrapper
->
pMgmt
,
pMsg
);
break
;
case
TDMT_DND_DROP_MNODE
:
case
TDMT_DND_DROP_QNODE
:
case
TDMT_DND_DROP_SNODE
:
case
TDMT_DND_DROP_BNODE
:
case
TDMT_DND_DROP_BNODE
:
pWrapper
=
dndGetWrapper
(
pDnode
,
BNODE
);
code
=
dndProcessDropNodeMsg
(
pMgmt
->
pDnode
,
pMsg
);
code
=
bmProcessDropReq
(
pWrapper
->
pMgmt
,
pMsg
);
break
;
case
TDMT_DND_CONFIG_DNODE
:
code
=
dmProcessConfigReq
(
pMgmt
,
pMsg
);
break
;
case
TDMT_MND_STATUS_RSP
:
code
=
dmProcessStatusRsp
(
pMgmt
,
pMsg
);
break
;
case
TDMT_MND_AUTH_RSP
:
code
=
dmProcessAuthRsp
(
pMgmt
,
pMsg
);
break
;
case
TDMT_MND_GRANT_RSP
:
code
=
dmProcessGrantRsp
(
pMgmt
,
pMsg
);
break
;
default:
default:
terrno
=
TSDB_CODE_MSG_NOT_PROCESSED
;
terrno
=
TSDB_CODE_MSG_NOT_PROCESSED
;
code
=
-
1
;
code
=
-
1
;
dError
(
"RPC %p, dnode msg:%s not processed"
,
pMsg
->
rpcMsg
.
handle
,
TMSG_INFO
(
msgType
));
dError
(
"RPC %p, dnode msg:%s not processed"
,
pMsg
->
rpcMsg
.
handle
,
TMSG_INFO
(
msgType
));
break
;
}
}
if
(
msgType
&
1u
)
{
if
(
msgType
&
1u
)
{
...
@@ -117,10 +82,9 @@ static void dmProcessQueue(SDnodeMgmt *pMgmt, SNodeMsg *pMsg) {
...
@@ -117,10 +82,9 @@ static void dmProcessQueue(SDnodeMgmt *pMgmt, SNodeMsg *pMsg) {
rpcSendResponse
(
&
rsp
);
rpcSendResponse
(
&
rsp
);
}
}
dTrace
(
"msg:%p, is freed"
,
pMsg
);
rpcFreeCont
(
pMsg
->
rpcMsg
.
pCont
);
rpcFreeCont
(
pMsg
->
rpcMsg
.
pCont
);
pMsg
->
rpcMsg
.
pCont
=
NULL
;
taosFreeQitem
(
pMsg
);
taosFreeQitem
(
pMsg
);
dTrace
(
"msg:%p, is freed"
,
pMsg
);
}
}
int32_t
dmStartWorker
(
SDnodeMgmt
*
pMgmt
)
{
int32_t
dmStartWorker
(
SDnodeMgmt
*
pMgmt
)
{
...
...
source/dnode/mgmt/test/sut/inc/server.h
浏览文件 @
43083fa3
...
@@ -28,7 +28,7 @@ class TestServer {
...
@@ -28,7 +28,7 @@ class TestServer {
private:
private:
SDnode
*
pDnode
;
SDnode
*
pDnode
;
pthread_t
*
threadId
;
pthread_t
threadId
;
char
path
[
PATH_MAX
];
char
path
[
PATH_MAX
];
char
fqdn
[
TSDB_FQDN_LEN
];
char
fqdn
[
TSDB_FQDN_LEN
];
char
firstEp
[
TSDB_EP_LEN
];
char
firstEp
[
TSDB_EP_LEN
];
...
...
source/dnode/mgmt/test/sut/src/server.cpp
浏览文件 @
43083fa3
...
@@ -16,10 +16,9 @@
...
@@ -16,10 +16,9 @@
#include "sut.h"
#include "sut.h"
void
*
serverLoop
(
void
*
param
)
{
void
*
serverLoop
(
void
*
param
)
{
while
(
1
)
{
SDnode
*
pDnode
=
(
SDnode
*
)
param
;
taosMsleep
(
100
);
dndRun
(
pDnode
);
pthread_testcancel
();
return
NULL
;
}
}
}
SDnodeOpt
TestServer
::
BuildOption
(
const
char
*
path
,
const
char
*
fqdn
,
uint16_t
port
,
const
char
*
firstEp
)
{
SDnodeOpt
TestServer
::
BuildOption
(
const
char
*
path
,
const
char
*
fqdn
,
uint16_t
port
,
const
char
*
firstEp
)
{
...
@@ -38,14 +37,16 @@ bool TestServer::DoStart() {
...
@@ -38,14 +37,16 @@ bool TestServer::DoStart() {
taosMkDir
(
path
);
taosMkDir
(
path
);
pDnode
=
dndCreate
(
&
option
);
pDnode
=
dndCreate
(
&
option
);
if
(
pDnode
!
=
NULL
)
{
if
(
pDnode
=
=
NULL
)
{
return
false
;
return
false
;
}
}
threadId
=
taosCreateThread
(
serverLoop
,
NULL
);
pthread_attr_t
thAttr
;
if
(
threadId
!=
NULL
)
{
pthread_attr_init
(
&
thAttr
);
return
false
;
pthread_attr_setdetachstate
(
&
thAttr
,
PTHREAD_CREATE_JOINABLE
);
}
pthread_create
(
&
threadId
,
&
thAttr
,
serverLoop
,
pDnode
);
pthread_attr_destroy
(
&
thAttr
);
taosMsleep
(
1000
);
return
true
;
return
true
;
}
}
...
@@ -67,10 +68,8 @@ bool TestServer::Start(const char* path, const char* fqdn, uint16_t port, const
...
@@ -67,10 +68,8 @@ bool TestServer::Start(const char* path, const char* fqdn, uint16_t port, const
}
}
void
TestServer
::
Stop
()
{
void
TestServer
::
Stop
()
{
if
(
threadId
!=
NULL
)
{
dndHandleEvent
(
pDnode
,
DND_EVENT_STOP
);
taosDestoryThread
(
threadId
);
pthread_join
(
threadId
,
NULL
);
threadId
=
NULL
;
}
if
(
pDnode
!=
NULL
)
{
if
(
pDnode
!=
NULL
)
{
dndClose
(
pDnode
);
dndClose
(
pDnode
);
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录