Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
慢慢CG
TDengine
提交
03da5549
T
TDengine
项目概览
慢慢CG
/
TDengine
与 Fork 源项目一致
Fork自
taosdata / TDengine
通知
1
Star
0
Fork
0
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
T
TDengine
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
提交
Issue看板
提交
03da5549
编写于
2月 24, 2020
作者:
S
slguan
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
get meta message
上级
41149af7
变更
17
隐藏空白更改
内联
并排
Showing
17 changed file
with
454 addition
and
532 deletion
+454
-532
src/client/inc/tsclient.h
src/client/inc/tsclient.h
+3
-3
src/client/src/tscParseInsert.c
src/client/src/tscParseInsert.c
+1
-1
src/client/src/tscSQLParser.c
src/client/src/tscSQLParser.c
+2
-2
src/client/src/tscServer.c
src/client/src/tscServer.c
+21
-21
src/client/src/tscUtil.c
src/client/src/tscUtil.c
+10
-10
src/inc/taosmsg.h
src/inc/taosmsg.h
+19
-18
src/mnode/inc/mgmtChildTable.h
src/mnode/inc/mgmtChildTable.h
+12
-8
src/mnode/inc/mgmtNormalTable.h
src/mnode/inc/mgmtNormalTable.h
+13
-9
src/mnode/inc/mgmtStreamTable.h
src/mnode/inc/mgmtStreamTable.h
+11
-7
src/mnode/inc/mgmtSuperTable.h
src/mnode/inc/mgmtSuperTable.h
+20
-12
src/mnode/inc/mgmtTable.h
src/mnode/inc/mgmtTable.h
+3
-5
src/mnode/src/mgmtChildTable.c
src/mnode/src/mgmtChildTable.c
+29
-2
src/mnode/src/mgmtNormalTable.c
src/mnode/src/mgmtNormalTable.c
+44
-2
src/mnode/src/mgmtShell.c
src/mnode/src/mgmtShell.c
+156
-423
src/mnode/src/mgmtStreamTable.c
src/mnode/src/mgmtStreamTable.c
+44
-3
src/mnode/src/mgmtSuperTable.c
src/mnode/src/mgmtSuperTable.c
+46
-4
src/mnode/src/mgmtTable.c
src/mnode/src/mgmtTable.c
+20
-2
未找到文件。
src/client/inc/tsclient.h
浏览文件 @
03da5549
...
@@ -120,7 +120,7 @@ typedef struct SCond {
...
@@ -120,7 +120,7 @@ typedef struct SCond {
}
SCond
;
}
SCond
;
typedef
struct
SJoinNode
{
typedef
struct
SJoinNode
{
char
meter
Id
[
TSDB_TABLE_ID_LEN
];
char
table
Id
[
TSDB_TABLE_ID_LEN
];
uint64_t
uid
;
uint64_t
uid
;
int16_t
tagCol
;
int16_t
tagCol
;
}
SJoinNode
;
}
SJoinNode
;
...
@@ -155,7 +155,7 @@ typedef struct SParamInfo {
...
@@ -155,7 +155,7 @@ typedef struct SParamInfo {
}
SParamInfo
;
}
SParamInfo
;
typedef
struct
STableDataBlocks
{
typedef
struct
STableDataBlocks
{
char
meter
Id
[
TSDB_TABLE_ID_LEN
];
char
table
Id
[
TSDB_TABLE_ID_LEN
];
int8_t
tsSource
;
// where does the UNIX timestamp come from, server or client
int8_t
tsSource
;
// where does the UNIX timestamp come from, server or client
bool
ordered
;
// if current rows are ordered or not
bool
ordered
;
// if current rows are ordered or not
int64_t
vgid
;
// virtual group id
int64_t
vgid
;
// virtual group id
...
@@ -384,7 +384,7 @@ int tscProcessSql(SSqlObj *pSql);
...
@@ -384,7 +384,7 @@ int tscProcessSql(SSqlObj *pSql);
void
tscAsyncInsertMultiVnodesProxy
(
void
*
param
,
TAOS_RES
*
tres
,
int
numOfRows
);
void
tscAsyncInsertMultiVnodesProxy
(
void
*
param
,
TAOS_RES
*
tres
,
int
numOfRows
);
int
tscRenewMeterMeta
(
SSqlObj
*
pSql
,
char
*
meter
Id
);
int
tscRenewMeterMeta
(
SSqlObj
*
pSql
,
char
*
table
Id
);
void
tscQueueAsyncRes
(
SSqlObj
*
pSql
);
void
tscQueueAsyncRes
(
SSqlObj
*
pSql
);
void
tscQueueAsyncError
(
void
(
*
fp
),
void
*
param
);
void
tscQueueAsyncError
(
void
(
*
fp
),
void
*
param
);
...
...
src/client/src/tscParseInsert.c
浏览文件 @
03da5549
...
@@ -1544,7 +1544,7 @@ void tscProcessMultiVnodesInsertFromFile(SSqlObj *pSql) {
...
@@ -1544,7 +1544,7 @@ void tscProcessMultiVnodesInsertFromFile(SSqlObj *pSql) {
continue
;
continue
;
}
}
strncpy
(
pMeterMetaInfo
->
name
,
pDataBlock
->
meter
Id
,
TSDB_TABLE_ID_LEN
);
strncpy
(
pMeterMetaInfo
->
name
,
pDataBlock
->
table
Id
,
TSDB_TABLE_ID_LEN
);
memset
(
pDataBlock
->
pData
,
0
,
pDataBlock
->
nAllocSize
);
memset
(
pDataBlock
->
pData
,
0
,
pDataBlock
->
nAllocSize
);
int32_t
ret
=
tscGetMeterMeta
(
pSql
,
pMeterMetaInfo
);
int32_t
ret
=
tscGetMeterMeta
(
pSql
,
pMeterMetaInfo
);
...
...
src/client/src/tscSQLParser.c
浏览文件 @
03da5549
...
@@ -2867,7 +2867,7 @@ static int32_t getJoinCondInfo(SQueryInfo* pQueryInfo, tSQLExpr* pExpr) {
...
@@ -2867,7 +2867,7 @@ static int32_t getJoinCondInfo(SQueryInfo* pQueryInfo, tSQLExpr* pExpr) {
pLeft
->
uid
=
pMeterMetaInfo
->
pMeterMeta
->
uid
;
pLeft
->
uid
=
pMeterMetaInfo
->
pMeterMeta
->
uid
;
pLeft
->
tagCol
=
tagColIndex
;
pLeft
->
tagCol
=
tagColIndex
;
strcpy
(
pLeft
->
meter
Id
,
pMeterMetaInfo
->
name
);
strcpy
(
pLeft
->
table
Id
,
pMeterMetaInfo
->
name
);
index
=
(
SColumnIndex
)
COLUMN_INDEX_INITIALIZER
;
index
=
(
SColumnIndex
)
COLUMN_INDEX_INITIALIZER
;
if
(
getColumnIndexByName
(
&
pExpr
->
pRight
->
colInfo
,
pQueryInfo
,
&
index
)
!=
TSDB_CODE_SUCCESS
)
{
if
(
getColumnIndexByName
(
&
pExpr
->
pRight
->
colInfo
,
pQueryInfo
,
&
index
)
!=
TSDB_CODE_SUCCESS
)
{
...
@@ -2879,7 +2879,7 @@ static int32_t getJoinCondInfo(SQueryInfo* pQueryInfo, tSQLExpr* pExpr) {
...
@@ -2879,7 +2879,7 @@ static int32_t getJoinCondInfo(SQueryInfo* pQueryInfo, tSQLExpr* pExpr) {
pRight
->
uid
=
pMeterMetaInfo
->
pMeterMeta
->
uid
;
pRight
->
uid
=
pMeterMetaInfo
->
pMeterMeta
->
uid
;
pRight
->
tagCol
=
tagColIndex
;
pRight
->
tagCol
=
tagColIndex
;
strcpy
(
pRight
->
meter
Id
,
pMeterMetaInfo
->
name
);
strcpy
(
pRight
->
table
Id
,
pMeterMetaInfo
->
name
);
pTagCond
->
joinInfo
.
hasJoin
=
true
;
pTagCond
->
joinInfo
.
hasJoin
=
true
;
return
TSDB_CODE_SUCCESS
;
return
TSDB_CODE_SUCCESS
;
...
...
src/client/src/tscServer.c
浏览文件 @
03da5549
...
@@ -1838,7 +1838,7 @@ int32_t tscBuildDropTableMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
...
@@ -1838,7 +1838,7 @@ int32_t tscBuildDropTableMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
pDropTableMsg
=
(
SDropTableMsg
*
)
pMsg
;
pDropTableMsg
=
(
SDropTableMsg
*
)
pMsg
;
strcpy
(
pDropTableMsg
->
meter
Id
,
pMeterMetaInfo
->
name
);
strcpy
(
pDropTableMsg
->
table
Id
,
pMeterMetaInfo
->
name
);
pDropTableMsg
->
igNotExists
=
pInfo
->
pDCLInfo
->
existsCheck
?
1
:
0
;
pDropTableMsg
->
igNotExists
=
pInfo
->
pDCLInfo
->
existsCheck
?
1
:
0
;
pMsg
+=
sizeof
(
SDropTableMsg
);
pMsg
+=
sizeof
(
SDropTableMsg
);
...
@@ -2353,7 +2353,7 @@ int tscBuildConnectMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
...
@@ -2353,7 +2353,7 @@ int tscBuildConnectMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
}
}
int
tscBuildMeterMetaMsg
(
SSqlObj
*
pSql
,
SSqlInfo
*
pInfo
)
{
int
tscBuildMeterMetaMsg
(
SSqlObj
*
pSql
,
SSqlInfo
*
pInfo
)
{
S
Meter
InfoMsg
*
pInfoMsg
;
S
Table
InfoMsg
*
pInfoMsg
;
char
*
pMsg
,
*
pStart
;
char
*
pMsg
,
*
pStart
;
int
msgLen
=
0
;
int
msgLen
=
0
;
...
@@ -2381,10 +2381,10 @@ int tscBuildMeterMetaMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
...
@@ -2381,10 +2381,10 @@ int tscBuildMeterMetaMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
pMsg
+=
sizeof
(
SMgmtHead
);
pMsg
+=
sizeof
(
SMgmtHead
);
pInfoMsg
=
(
S
Meter
InfoMsg
*
)
pMsg
;
pInfoMsg
=
(
S
Table
InfoMsg
*
)
pMsg
;
strcpy
(
pInfoMsg
->
meter
Id
,
pMeterMetaInfo
->
name
);
strcpy
(
pInfoMsg
->
table
Id
,
pMeterMetaInfo
->
name
);
pInfoMsg
->
createFlag
=
htons
(
pSql
->
cmd
.
createOnDemand
?
1
:
0
);
pInfoMsg
->
createFlag
=
htons
(
pSql
->
cmd
.
createOnDemand
?
1
:
0
);
pMsg
+=
sizeof
(
S
Meter
InfoMsg
);
pMsg
+=
sizeof
(
S
Table
InfoMsg
);
if
(
pSql
->
cmd
.
createOnDemand
)
{
if
(
pSql
->
cmd
.
createOnDemand
)
{
memcpy
(
pInfoMsg
->
tags
,
tmpData
,
sizeof
(
STagData
));
memcpy
(
pInfoMsg
->
tags
,
tmpData
,
sizeof
(
STagData
));
...
@@ -2403,7 +2403,7 @@ int tscBuildMeterMetaMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
...
@@ -2403,7 +2403,7 @@ int tscBuildMeterMetaMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
/**
/**
* multi meter meta req pkg format:
* multi meter meta req pkg format:
* | SMgmtHead | SMulti
MeterInfoMsg | meterId0 | meterId1 | meter
Id2 | ......
* | SMgmtHead | SMulti
TableInfoMsg | tableId0 | tableId1 | table
Id2 | ......
* no used 4B
* no used 4B
**/
**/
int
tscBuildMultiMeterMetaMsg
(
SSqlObj
*
pSql
,
SSqlInfo
*
pInfo
)
{
int
tscBuildMultiMeterMetaMsg
(
SSqlObj
*
pSql
,
SSqlInfo
*
pInfo
)
{
...
@@ -2421,16 +2421,16 @@ int tscBuildMultiMeterMetaMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
...
@@ -2421,16 +2421,16 @@ int tscBuildMultiMeterMetaMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
SMgmtHead
*
pMgmt
=
(
SMgmtHead
*
)(
pCmd
->
payload
+
tsRpcHeadSize
);
SMgmtHead
*
pMgmt
=
(
SMgmtHead
*
)(
pCmd
->
payload
+
tsRpcHeadSize
);
memset
(
pMgmt
->
db
,
0
,
TSDB_TABLE_ID_LEN
);
// server don't need the db
memset
(
pMgmt
->
db
,
0
,
TSDB_TABLE_ID_LEN
);
// server don't need the db
SMulti
MeterInfoMsg
*
pInfoMsg
=
(
SMultiMeter
InfoMsg
*
)(
pCmd
->
payload
+
tsRpcHeadSize
+
sizeof
(
SMgmtHead
));
SMulti
TableInfoMsg
*
pInfoMsg
=
(
SMultiTable
InfoMsg
*
)(
pCmd
->
payload
+
tsRpcHeadSize
+
sizeof
(
SMgmtHead
));
pInfoMsg
->
numOf
Meter
s
=
htonl
((
int32_t
)
pCmd
->
count
);
pInfoMsg
->
numOf
Table
s
=
htonl
((
int32_t
)
pCmd
->
count
);
if
(
pCmd
->
payloadLen
>
0
)
{
if
(
pCmd
->
payloadLen
>
0
)
{
memcpy
(
pInfoMsg
->
meterId
,
tmpData
,
pCmd
->
payloadLen
);
memcpy
(
pInfoMsg
->
tableIds
,
tmpData
,
pCmd
->
payloadLen
);
}
}
tfree
(
tmpData
);
tfree
(
tmpData
);
pCmd
->
payloadLen
+=
sizeof
(
SMgmtHead
)
+
sizeof
(
SMulti
Meter
InfoMsg
);
pCmd
->
payloadLen
+=
sizeof
(
SMgmtHead
)
+
sizeof
(
SMulti
Table
InfoMsg
);
pCmd
->
msgType
=
TSDB_MSG_TYPE_MULTI_TABLE_META
;
pCmd
->
msgType
=
TSDB_MSG_TYPE_MULTI_TABLE_META
;
assert
(
pCmd
->
payloadLen
+
minMsgSize
()
<=
pCmd
->
allocSize
);
assert
(
pCmd
->
payloadLen
+
minMsgSize
()
<=
pCmd
->
allocSize
);
...
@@ -2502,13 +2502,13 @@ int tscBuildMetricMetaMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
...
@@ -2502,13 +2502,13 @@ int tscBuildMetricMetaMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
// todo refactor
// todo refactor
pMetaMsg
->
joinCondLen
=
htonl
((
TSDB_TABLE_ID_LEN
+
sizeof
(
int16_t
))
*
2
);
pMetaMsg
->
joinCondLen
=
htonl
((
TSDB_TABLE_ID_LEN
+
sizeof
(
int16_t
))
*
2
);
memcpy
(
pMsg
,
pTagCond
->
joinInfo
.
left
.
meter
Id
,
TSDB_TABLE_ID_LEN
);
memcpy
(
pMsg
,
pTagCond
->
joinInfo
.
left
.
table
Id
,
TSDB_TABLE_ID_LEN
);
pMsg
+=
TSDB_TABLE_ID_LEN
;
pMsg
+=
TSDB_TABLE_ID_LEN
;
*
(
int16_t
*
)
pMsg
=
pTagCond
->
joinInfo
.
left
.
tagCol
;
*
(
int16_t
*
)
pMsg
=
pTagCond
->
joinInfo
.
left
.
tagCol
;
pMsg
+=
sizeof
(
int16_t
);
pMsg
+=
sizeof
(
int16_t
);
memcpy
(
pMsg
,
pTagCond
->
joinInfo
.
right
.
meter
Id
,
TSDB_TABLE_ID_LEN
);
memcpy
(
pMsg
,
pTagCond
->
joinInfo
.
right
.
table
Id
,
TSDB_TABLE_ID_LEN
);
pMsg
+=
TSDB_TABLE_ID_LEN
;
pMsg
+=
TSDB_TABLE_ID_LEN
;
*
(
int16_t
*
)
pMsg
=
pTagCond
->
joinInfo
.
right
.
tagCol
;
*
(
int16_t
*
)
pMsg
=
pTagCond
->
joinInfo
.
right
.
tagCol
;
...
@@ -2590,7 +2590,7 @@ int tscBuildMetricMetaMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
...
@@ -2590,7 +2590,7 @@ int tscBuildMetricMetaMsg(SSqlObj *pSql, SSqlInfo *pInfo) {
}
}
}
}
strcpy
(
pElem
->
meter
Id
,
pMeterMetaInfo
->
name
);
strcpy
(
pElem
->
table
Id
,
pMeterMetaInfo
->
name
);
pElem
->
numOfTags
=
htons
(
pMeterMetaInfo
->
numOfTags
);
pElem
->
numOfTags
=
htons
(
pMeterMetaInfo
->
numOfTags
);
int16_t
len
=
pMsg
-
(
char
*
)
pElem
;
int16_t
len
=
pMsg
-
(
char
*
)
pElem
;
...
@@ -2751,7 +2751,7 @@ int tscProcessMeterMetaRsp(SSqlObj *pSql) {
...
@@ -2751,7 +2751,7 @@ int tscProcessMeterMetaRsp(SSqlObj *pSql) {
/**
/**
* multi meter meta rsp pkg format:
* multi meter meta rsp pkg format:
* | STaosRsp | ieType | SMulti
Meter
InfoMsg | SMeterMeta0 | SSchema0 | SMeterMeta1 | SSchema1 | SMeterMeta2 | SSchema2
* | STaosRsp | ieType | SMulti
Table
InfoMsg | SMeterMeta0 | SSchema0 | SMeterMeta1 | SSchema1 | SMeterMeta2 | SSchema2
* |...... 1B 1B 4B
* |...... 1B 1B 4B
**/
**/
int
tscProcessMultiMeterMetaRsp
(
SSqlObj
*
pSql
)
{
int
tscProcessMultiMeterMetaRsp
(
SSqlObj
*
pSql
)
{
...
@@ -2772,13 +2772,13 @@ int tscProcessMultiMeterMetaRsp(SSqlObj *pSql) {
...
@@ -2772,13 +2772,13 @@ int tscProcessMultiMeterMetaRsp(SSqlObj *pSql) {
rsp
++
;
rsp
++
;
SMulti
MeterInfoMsg
*
pInfo
=
(
SMultiMeter
InfoMsg
*
)
rsp
;
SMulti
TableInfoMsg
*
pInfo
=
(
SMultiTable
InfoMsg
*
)
rsp
;
totalNum
=
htonl
(
pInfo
->
numOf
Meter
s
);
totalNum
=
htonl
(
pInfo
->
numOf
Table
s
);
rsp
+=
sizeof
(
SMulti
Meter
InfoMsg
);
rsp
+=
sizeof
(
SMulti
Table
InfoMsg
);
for
(
i
=
0
;
i
<
totalNum
;
i
++
)
{
for
(
i
=
0
;
i
<
totalNum
;
i
++
)
{
SMultiMeterMeta
*
pMultiMeta
=
(
SMultiMeterMeta
*
)
rsp
;
SMultiMeterMeta
*
pMultiMeta
=
(
SMultiMeterMeta
*
)
rsp
;
SMeterMeta
*
pMeta
=
&
pMultiMeta
->
meta
;
SMeterMeta
*
pMeta
=
&
pMultiMeta
->
meta
s
;
pMeta
->
sid
=
htonl
(
pMeta
->
sid
);
pMeta
->
sid
=
htonl
(
pMeta
->
sid
);
pMeta
->
sversion
=
htons
(
pMeta
->
sversion
);
pMeta
->
sversion
=
htons
(
pMeta
->
sversion
);
...
@@ -2850,7 +2850,7 @@ int tscProcessMultiMeterMetaRsp(SSqlObj *pSql) {
...
@@ -2850,7 +2850,7 @@ int tscProcessMultiMeterMetaRsp(SSqlObj *pSql) {
int32_t
size
=
(
int32_t
)(
rsp
-
((
char
*
)
pMeta
));
// Consistent with SMeterMeta in cache
int32_t
size
=
(
int32_t
)(
rsp
-
((
char
*
)
pMeta
));
// Consistent with SMeterMeta in cache
pMeta
->
index
=
0
;
pMeta
->
index
=
0
;
(
void
)
taosAddDataIntoCache
(
tscCacheHandle
,
pM
ultiMeta
->
meter
Id
,
(
char
*
)
pMeta
,
size
,
tsMeterMetaKeepTimer
);
(
void
)
taosAddDataIntoCache
(
tscCacheHandle
,
pM
eta
->
table
Id
,
(
char
*
)
pMeta
,
size
,
tsMeterMetaKeepTimer
);
}
}
pSql
->
res
.
code
=
TSDB_CODE_SUCCESS
;
pSql
->
res
.
code
=
TSDB_CODE_SUCCESS
;
...
@@ -3312,10 +3312,10 @@ static void tscWaitingForCreateTable(SSqlCmd *pCmd) {
...
@@ -3312,10 +3312,10 @@ static void tscWaitingForCreateTable(SSqlCmd *pCmd) {
/**
/**
* in renew metermeta, do not retrieve metadata in cache.
* in renew metermeta, do not retrieve metadata in cache.
* @param pSql sql object
* @param pSql sql object
* @param
meter
Id meter id
* @param
table
Id meter id
* @return status code
* @return status code
*/
*/
int
tscRenewMeterMeta
(
SSqlObj
*
pSql
,
char
*
meter
Id
)
{
int
tscRenewMeterMeta
(
SSqlObj
*
pSql
,
char
*
table
Id
)
{
int
code
=
0
;
int
code
=
0
;
// handle metric meta renew process
// handle metric meta renew process
...
...
src/client/src/tscUtil.c
浏览文件 @
03da5549
...
@@ -56,7 +56,7 @@ void tscGetMetricMetaCacheKey(SQueryInfo* pQueryInfo, char* str, uint64_t uid) {
...
@@ -56,7 +56,7 @@ void tscGetMetricMetaCacheKey(SQueryInfo* pQueryInfo, char* str, uint64_t uid) {
char
join
[
512
]
=
{
0
};
char
join
[
512
]
=
{
0
};
if
(
pTagCond
->
joinInfo
.
hasJoin
)
{
if
(
pTagCond
->
joinInfo
.
hasJoin
)
{
sprintf
(
join
,
"%s,%s"
,
pTagCond
->
joinInfo
.
left
.
meterId
,
pTagCond
->
joinInfo
.
right
.
meter
Id
);
sprintf
(
join
,
"%s,%s"
,
pTagCond
->
joinInfo
.
left
.
tableId
,
pTagCond
->
joinInfo
.
right
.
table
Id
);
}
}
// estimate the buffer size
// estimate the buffer size
...
@@ -156,13 +156,13 @@ bool tscIsSelectivityWithTagQuery(SSqlCmd* pCmd) {
...
@@ -156,13 +156,13 @@ bool tscIsSelectivityWithTagQuery(SSqlCmd* pCmd) {
return
false
;
return
false
;
}
}
void
tscGetDBInfoFromMeterId
(
char
*
meter
Id
,
char
*
db
)
{
void
tscGetDBInfoFromMeterId
(
char
*
table
Id
,
char
*
db
)
{
char
*
st
=
strstr
(
meter
Id
,
TS_PATH_DELIMITER
);
char
*
st
=
strstr
(
table
Id
,
TS_PATH_DELIMITER
);
if
(
st
!=
NULL
)
{
if
(
st
!=
NULL
)
{
char
*
end
=
strstr
(
st
+
1
,
TS_PATH_DELIMITER
);
char
*
end
=
strstr
(
st
+
1
,
TS_PATH_DELIMITER
);
if
(
end
!=
NULL
)
{
if
(
end
!=
NULL
)
{
memcpy
(
db
,
meterId
,
(
end
-
meter
Id
));
memcpy
(
db
,
tableId
,
(
end
-
table
Id
));
db
[
end
-
meter
Id
]
=
0
;
db
[
end
-
table
Id
]
=
0
;
return
;
return
;
}
}
}
}
...
@@ -590,12 +590,12 @@ int32_t tscCopyDataBlockToPayload(SSqlObj* pSql, STableDataBlocks* pDataBlock) {
...
@@ -590,12 +590,12 @@ int32_t tscCopyDataBlockToPayload(SSqlObj* pSql, STableDataBlocks* pDataBlock) {
// set the correct metermeta object, the metermeta has been locked in pDataBlocks, so it must be in the cache
// set the correct metermeta object, the metermeta has been locked in pDataBlocks, so it must be in the cache
if
(
pMeterMetaInfo
->
pMeterMeta
!=
pDataBlock
->
pMeterMeta
)
{
if
(
pMeterMetaInfo
->
pMeterMeta
!=
pDataBlock
->
pMeterMeta
)
{
strcpy
(
pMeterMetaInfo
->
name
,
pDataBlock
->
meter
Id
);
strcpy
(
pMeterMetaInfo
->
name
,
pDataBlock
->
table
Id
);
taosRemoveDataFromCache
(
tscCacheHandle
,
(
void
**
)
&
(
pMeterMetaInfo
->
pMeterMeta
),
false
);
taosRemoveDataFromCache
(
tscCacheHandle
,
(
void
**
)
&
(
pMeterMetaInfo
->
pMeterMeta
),
false
);
pMeterMetaInfo
->
pMeterMeta
=
taosTransferDataInCache
(
tscCacheHandle
,
(
void
**
)
&
pDataBlock
->
pMeterMeta
);
pMeterMetaInfo
->
pMeterMeta
=
taosTransferDataInCache
(
tscCacheHandle
,
(
void
**
)
&
pDataBlock
->
pMeterMeta
);
}
else
{
}
else
{
assert
(
strncmp
(
pMeterMetaInfo
->
name
,
pDataBlock
->
meterId
,
tListLen
(
pDataBlock
->
meter
Id
))
==
0
);
assert
(
strncmp
(
pMeterMetaInfo
->
name
,
pDataBlock
->
tableId
,
tListLen
(
pDataBlock
->
table
Id
))
==
0
);
}
}
/*
/*
...
@@ -660,7 +660,7 @@ int32_t tscCreateDataBlock(size_t initialSize, int32_t rowSize, int32_t startOff
...
@@ -660,7 +660,7 @@ int32_t tscCreateDataBlock(size_t initialSize, int32_t rowSize, int32_t startOff
dataBuf
->
size
=
startOffset
;
dataBuf
->
size
=
startOffset
;
dataBuf
->
tsSource
=
-
1
;
dataBuf
->
tsSource
=
-
1
;
strncpy
(
dataBuf
->
meter
Id
,
name
,
TSDB_TABLE_ID_LEN
);
strncpy
(
dataBuf
->
table
Id
,
name
,
TSDB_TABLE_ID_LEN
);
/*
/*
* The metermeta may be released since the metermeta cache are completed clean by other thread
* The metermeta may be released since the metermeta cache are completed clean by other thread
...
@@ -709,7 +709,7 @@ int32_t tscMergeTableDataBlocks(SSqlObj* pSql, SDataBlockList* pTableDataBlockLi
...
@@ -709,7 +709,7 @@ int32_t tscMergeTableDataBlocks(SSqlObj* pSql, SDataBlockList* pTableDataBlockLi
STableDataBlocks
*
dataBuf
=
NULL
;
STableDataBlocks
*
dataBuf
=
NULL
;
int32_t
ret
=
int32_t
ret
=
tscGetDataBlockFromList
(
pVnodeDataBlockHashList
,
pVnodeDataBlockList
,
pOneTableBlock
->
vgid
,
TSDB_PAYLOAD_SIZE
,
tscGetDataBlockFromList
(
pVnodeDataBlockHashList
,
pVnodeDataBlockList
,
pOneTableBlock
->
vgid
,
TSDB_PAYLOAD_SIZE
,
tsInsertHeadSize
,
0
,
pOneTableBlock
->
meter
Id
,
pOneTableBlock
->
pMeterMeta
,
&
dataBuf
);
tsInsertHeadSize
,
0
,
pOneTableBlock
->
table
Id
,
pOneTableBlock
->
pMeterMeta
,
&
dataBuf
);
if
(
ret
!=
TSDB_CODE_SUCCESS
)
{
if
(
ret
!=
TSDB_CODE_SUCCESS
)
{
tscError
(
"%p failed to prepare the data block buffer for merging table data, code:%d"
,
pSql
,
ret
);
tscError
(
"%p failed to prepare the data block buffer for merging table data, code:%d"
,
pSql
,
ret
);
taosCleanUpHashTable
(
pVnodeDataBlockHashList
);
taosCleanUpHashTable
(
pVnodeDataBlockHashList
);
...
@@ -743,7 +743,7 @@ int32_t tscMergeTableDataBlocks(SSqlObj* pSql, SDataBlockList* pTableDataBlockLi
...
@@ -743,7 +743,7 @@ int32_t tscMergeTableDataBlocks(SSqlObj* pSql, SDataBlockList* pTableDataBlockLi
char
*
e
=
(
char
*
)
pBlocks
->
payLoad
+
pOneTableBlock
->
rowSize
*
(
pBlocks
->
numOfRows
-
1
);
char
*
e
=
(
char
*
)
pBlocks
->
payLoad
+
pOneTableBlock
->
rowSize
*
(
pBlocks
->
numOfRows
-
1
);
tscTrace
(
"%p
meterId:%s, sid:%d rows:%d sversion:%d skey:%"
PRId64
", ekey:%"
PRId64
,
pSql
,
pOneTableBlock
->
meter
Id
,
pBlocks
->
sid
,
tscTrace
(
"%p
tableId:%s, sid:%d rows:%d sversion:%d skey:%"
PRId64
", ekey:%"
PRId64
,
pSql
,
pOneTableBlock
->
table
Id
,
pBlocks
->
sid
,
pBlocks
->
numOfRows
,
pBlocks
->
sversion
,
GET_INT64_VAL
(
pBlocks
->
payLoad
),
GET_INT64_VAL
(
e
));
pBlocks
->
numOfRows
,
pBlocks
->
sversion
,
GET_INT64_VAL
(
pBlocks
->
payLoad
),
GET_INT64_VAL
(
e
));
pBlocks
->
sid
=
htonl
(
pBlocks
->
sid
);
pBlocks
->
sid
=
htonl
(
pBlocks
->
sid
);
...
...
src/inc/taosmsg.h
浏览文件 @
03da5549
...
@@ -277,11 +277,11 @@ typedef struct {
...
@@ -277,11 +277,11 @@ typedef struct {
int16_t
numOfColumns
;
int16_t
numOfColumns
;
int16_t
sqlLen
;
// the length of SQL, it starts after schema , sql is a null-terminated string
int16_t
sqlLen
;
// the length of SQL, it starts after schema , sql is a null-terminated string
int16_t
reserved
[
16
];
int16_t
reserved
[
16
];
SSchema
schema
[];
SSchema
schema
[];
}
SCreateTableMsg
;
}
SCreateTableMsg
;
typedef
struct
{
typedef
struct
{
char
meter
Id
[
TSDB_TABLE_ID_LEN
];
char
table
Id
[
TSDB_TABLE_ID_LEN
];
char
db
[
TSDB_DB_NAME_LEN
];
char
db
[
TSDB_DB_NAME_LEN
];
int8_t
igNotExists
;
int8_t
igNotExists
;
}
SDropTableMsg
;
}
SDropTableMsg
;
...
@@ -348,7 +348,7 @@ typedef struct {
...
@@ -348,7 +348,7 @@ typedef struct {
short
vnode
;
short
vnode
;
int32_t
sid
;
int32_t
sid
;
uint64_t
uid
;
uint64_t
uid
;
char
meter
Id
[
TSDB_TABLE_ID_LEN
];
char
table
Id
[
TSDB_TABLE_ID_LEN
];
}
SDRemoveTableMsg
;
}
SDRemoveTableMsg
;
typedef
struct
{
typedef
struct
{
...
@@ -615,7 +615,7 @@ typedef struct {
...
@@ -615,7 +615,7 @@ typedef struct {
typedef
struct
{
typedef
struct
{
uint32_t
destId
;
uint32_t
destId
;
uint32_t
destIp
;
uint32_t
destIp
;
char
meter
Id
[
TSDB_UNI_LEN
];
char
table
Id
[
TSDB_UNI_LEN
];
char
empty
[
3
];
char
empty
[
3
];
uint8_t
msgType
;
uint8_t
msgType
;
int32_t
msgLen
;
int32_t
msgLen
;
...
@@ -647,20 +647,21 @@ typedef struct {
...
@@ -647,20 +647,21 @@ typedef struct {
}
SVPeersMsg
;
}
SVPeersMsg
;
typedef
struct
{
typedef
struct
{
char
meterId
[
TSDB_TABLE_ID_LEN
];
char
tableId
[
TSDB_TABLE_ID_LEN
];
short
createFlag
;
char
db
[
TSDB_DB_NAME_LEN
];
char
tags
[];
int16_t
createFlag
;
}
SMeterInfoMsg
;
char
tags
[];
}
STableInfoMsg
;
typedef
struct
{
typedef
struct
{
int32_t
numOf
Meter
s
;
int32_t
numOf
Table
s
;
char
meterId
[];
char
tableIds
[];
}
SMulti
Meter
InfoMsg
;
}
SMulti
Table
InfoMsg
;
typedef
struct
{
typedef
struct
{
int16_t
elemLen
;
int16_t
elemLen
;
char
meter
Id
[
TSDB_TABLE_ID_LEN
];
char
table
Id
[
TSDB_TABLE_ID_LEN
];
int16_t
orderIndex
;
int16_t
orderIndex
;
int16_t
orderType
;
// used in group by xx order by xxx
int16_t
orderType
;
// used in group by xx order by xxx
...
@@ -701,26 +702,26 @@ typedef struct {
...
@@ -701,26 +702,26 @@ typedef struct {
}
SMetricMeta
;
}
SMetricMeta
;
typedef
struct
SMeterMeta
{
typedef
struct
SMeterMeta
{
char
tableId
[
TSDB_TABLE_ID_LEN
];
// note: This field must be at the front
int32_t
contLen
;
uint8_t
numOfTags
:
6
;
uint8_t
numOfTags
:
6
;
uint8_t
precision
:
2
;
uint8_t
precision
:
2
;
uint8_t
tableType
:
4
;
uint8_t
tableType
:
4
;
uint8_t
index
:
4
;
// used locally
uint8_t
index
:
4
;
// used locally
int16_t
numOfColumns
;
int16_t
numOfColumns
;
int16_t
rowSize
;
// used locally, calculated in client
int16_t
rowSize
;
// used locally, calculated in client
int16_t
sversion
;
int16_t
sversion
;
SVPeerDesc
vpeerDesc
[
TSDB_VNODES_SUPPORT
];
SVPeerDesc
vpeerDesc
[
TSDB_VNODES_SUPPORT
];
int32_t
sid
;
int32_t
sid
;
int32_t
vgid
;
int32_t
vgid
;
uint64_t
uid
;
uint64_t
uid
;
SSchema
schema
[];
}
SMeterMeta
;
}
SMeterMeta
;
typedef
struct
SMultiMeterMeta
{
typedef
struct
SMultiMeterMeta
{
char
meterId
[
TSDB_TABLE_ID_LEN
];
// note: This field must be at the front
int32_t
numOfTables
;
SMeterMeta
meta
;
int32_t
contLen
;
SMeterMeta
metas
[];
}
SMultiMeterMeta
;
}
SMultiMeterMeta
;
typedef
struct
{
typedef
struct
{
...
...
src/mnode/inc/mgmtChildTable.h
浏览文件 @
03da5549
...
@@ -26,14 +26,18 @@ extern "C" {
...
@@ -26,14 +26,18 @@ extern "C" {
#include "mnode.h"
#include "mnode.h"
int32_t
mgmtInitChildTables
();
int32_t
mgmtInitChildTables
();
void
mgmtCleanUpChildTables
();
void
mgmtCleanUpChildTables
();
int32_t
mgmtCreateChildTable
(
SDbObj
*
pDb
,
SCreateTableMsg
*
pCreate
,
SVgObj
*
pVgroup
,
int32_t
sid
);
int32_t
mgmtDropChildTable
(
SDbObj
*
pDb
,
SChildTableObj
*
pTable
);
void
*
mgmtGetChildTable
(
char
*
tableId
);
int32_t
mgmtAlterChildTable
(
SDbObj
*
pDb
,
SAlterTableMsg
*
pAlter
);
int32_t
mgmtModifyChildTableTagValueByName
(
SChildTableObj
*
pTable
,
char
*
tagName
,
char
*
nContent
);
int32_t
mgmtCreateChildTable
(
SDbObj
*
pDb
,
SCreateTableMsg
*
pCreate
,
SVgObj
*
pVgroup
,
int32_t
sid
);
SChildTableObj
*
mgmtGetChildTable
(
char
*
tableId
);
int32_t
mgmtDropChildTable
(
SDbObj
*
pDb
,
SChildTableObj
*
pTable
);
int8_t
*
mgmtBuildCreateChildTableMsg
(
SChildTableObj
*
pTable
,
SVgObj
*
pVgroup
);
int32_t
mgmtAlterChildTable
(
SDbObj
*
pDb
,
SAlterTableMsg
*
pAlter
);
int32_t
mgmtModifyChildTableTagValueByName
(
SChildTableObj
*
pTable
,
char
*
tagName
,
char
*
nContent
);
int8_t
*
mgmtBuildCreateChildTableMsg
(
SChildTableObj
*
pTable
,
SVgObj
*
pVgroup
);
int32_t
mgmtGetChildTableMeta
(
SDbObj
*
pDb
,
SChildTableObj
*
pTable
,
SMeterMeta
*
pMeta
,
bool
usePublicIp
);
#ifdef __cplusplus
#ifdef __cplusplus
}
}
...
...
src/mnode/inc/mgmtNormalTable.h
浏览文件 @
03da5549
...
@@ -23,15 +23,19 @@ extern "C" {
...
@@ -23,15 +23,19 @@ extern "C" {
#include <stdint.h>
#include <stdint.h>
#include <stdbool.h>
#include <stdbool.h>
#include "mnode.h"
#include "mnode.h"
int32_t
mgmtInitNormalTables
();
int32_t
mgmtInitNormalTables
();
void
mgmtCleanUpNormalTables
();
void
mgmtCleanUpNormalTables
();
int32_t
mgmtCreateNormalTable
(
SDbObj
*
pDb
,
SCreateTableMsg
*
pCreate
,
SVgObj
*
pVgroup
,
int32_t
sid
);
int32_t
mgmtDropNormalTable
(
SDbObj
*
pDb
,
SNormalTableObj
*
pTable
);
void
*
mgmtGetNormalTable
(
char
*
tableId
);
int32_t
mgmtAddNormalTableColumn
(
SNormalTableObj
*
pTable
,
SSchema
schema
[],
int32_t
ncols
);
int32_t
mgmtDropNormalTableColumnByName
(
SNormalTableObj
*
pTable
,
char
*
colName
);
int32_t
mgmtCreateNormalTable
(
SDbObj
*
pDb
,
SCreateTableMsg
*
pCreate
,
SVgObj
*
pVgroup
,
int32_t
sid
);
SNormalTableObj
*
mgmtGetNormalTable
(
char
*
tableId
);
int32_t
mgmtDropNormalTable
(
SDbObj
*
pDb
,
SNormalTableObj
*
pTable
);
int8_t
*
mgmtBuildCreateNormalTableMsg
(
SNormalTableObj
*
pTable
);
int32_t
mgmtAddNormalTableColumn
(
SNormalTableObj
*
pTable
,
SSchema
schema
[],
int32_t
ncols
);
int32_t
mgmtDropNormalTableColumnByName
(
SNormalTableObj
*
pTable
,
char
*
colName
);
int8_t
*
mgmtBuildCreateNormalTableMsg
(
SNormalTableObj
*
pTable
);
int32_t
mgmtGetNormalTableMeta
(
SDbObj
*
pDb
,
SNormalTableObj
*
pTable
,
SMeterMeta
*
pMeta
,
bool
usePublicIp
);
#ifdef __cplusplus
#ifdef __cplusplus
}
}
...
...
src/mnode/inc/mgmtStreamTable.h
浏览文件 @
03da5549
...
@@ -24,13 +24,17 @@ extern "C" {
...
@@ -24,13 +24,17 @@ extern "C" {
#include <stdbool.h>
#include <stdbool.h>
#include "mnode.h"
#include "mnode.h"
int32_t
mgmtInitStreamTables
();
int32_t
mgmtInitStreamTables
();
void
mgmtCleanUpStreamTables
();
void
mgmtCleanUpStreamTables
();
int32_t
mgmtCreateStreamTable
(
SDbObj
*
pDb
,
SCreateTableMsg
*
pCreate
,
SVgObj
*
pVgroup
,
int32_t
sid
);
int32_t
mgmtDropStreamTable
(
SDbObj
*
pDb
,
SStreamTableObj
*
pTable
);
void
*
mgmtGetStreamTable
(
char
*
tableId
);
int32_t
mgmtAlterStreamTable
(
SDbObj
*
pDb
,
SAlterTableMsg
*
pAlter
);
SStreamTableObj
*
mgmtGetStreamTable
(
char
*
tableId
);
int32_t
mgmtCreateStreamTable
(
SDbObj
*
pDb
,
SCreateTableMsg
*
pCreate
,
SVgObj
*
pVgroup
,
int32_t
sid
);
int8_t
*
mgmtBuildCreateStreamTableMsg
(
SStreamTableObj
*
pTable
,
SVgObj
*
pVgroup
);
int32_t
mgmtDropStreamTable
(
SDbObj
*
pDb
,
SStreamTableObj
*
pTable
);
int32_t
mgmtAlterStreamTable
(
SDbObj
*
pDb
,
SAlterTableMsg
*
pAlter
);
int8_t
*
mgmtBuildCreateStreamTableMsg
(
SStreamTableObj
*
pTable
,
SVgObj
*
pVgroup
);
int32_t
mgmtGetStreamTableMeta
(
SDbObj
*
pDb
,
SStreamTableObj
*
pTable
,
SMeterMeta
*
pMeta
,
bool
usePublicIp
);
#ifdef __cplusplus
#ifdef __cplusplus
}
}
...
...
src/mnode/inc/mgmtSuperTable.h
浏览文件 @
03da5549
...
@@ -26,18 +26,26 @@ extern "C" {
...
@@ -26,18 +26,26 @@ extern "C" {
#include "taosdef.h"
#include "taosdef.h"
#include "mnode.h"
#include "mnode.h"
int32_t
mgmtInitSuperTables
();
int32_t
mgmtInitSuperTables
();
void
mgmtCleanUpSuperTables
();
void
mgmtCleanUpSuperTables
();
int32_t
mgmtCreateSuperTable
(
SDbObj
*
pDb
,
SCreateTableMsg
*
pCreate
);
int32_t
mgmtDropSuperTable
(
SDbObj
*
pDb
,
SSuperTableObj
*
pTable
);
void
*
mgmtGetSuperTable
(
char
*
tableId
);
SSuperTableObj
*
mgmtGetSuperTable
(
char
*
tableId
);
int32_t
mgmtGetShowSuperTableMeta
(
SMeterMeta
*
pMeta
,
SShowObj
*
pShow
,
void
*
pConn
);
int32_t
mgmtFindSuperTableTagIndex
(
SSuperTableObj
*
pTable
,
const
char
*
tagName
);
int32_t
mgmtRetrieveShowSuperTables
(
SShowObj
*
pShow
,
char
*
data
,
int32_t
rows
,
void
*
pConn
);
int32_t
mgmtAddSuperTableTag
(
SSuperTableObj
*
pTable
,
SSchema
schema
[],
int32_t
ntags
);
int32_t
mgmtDropSuperTableTag
(
SSuperTableObj
*
pTable
,
char
*
tagName
);
int32_t
mgmtCreateSuperTable
(
SDbObj
*
pDb
,
SCreateTableMsg
*
pCreate
);
int32_t
mgmtModifySuperTableTagNameByName
(
SSuperTableObj
*
pTable
,
char
*
oldTagName
,
char
*
newTagName
);
int32_t
mgmtDropSuperTable
(
SDbObj
*
pDb
,
SSuperTableObj
*
pTable
);
int32_t
mgmtAddSuperTableColumn
(
SSuperTableObj
*
pTable
,
SSchema
schema
[],
int32_t
ncols
);
int32_t
mgmtAddSuperTableTag
(
SSuperTableObj
*
pTable
,
SSchema
schema
[],
int32_t
ntags
);
int32_t
mgmtDropSuperTableColumnByName
(
SSuperTableObj
*
pTable
,
char
*
colName
);
int32_t
mgmtDropSuperTableTag
(
SSuperTableObj
*
pTable
,
char
*
tagName
);
int32_t
mgmtGetTagsLength
(
SSuperTableObj
*
pSuperTable
,
int32_t
col
);
int32_t
mgmtModifySuperTableTagNameByName
(
SSuperTableObj
*
pTable
,
char
*
oldTagName
,
char
*
newTagName
);
int32_t
mgmtAddSuperTableColumn
(
SSuperTableObj
*
pTable
,
SSchema
schema
[],
int32_t
ncols
);
int32_t
mgmtDropSuperTableColumnByName
(
SSuperTableObj
*
pTable
,
char
*
colName
);
int32_t
mgmtGetSuperTableMeta
(
SDbObj
*
pDb
,
SSuperTableObj
*
pTable
,
SMeterMeta
*
pMeta
,
bool
usePublicIp
);
int32_t
mgmtFindSuperTableTagIndex
(
SSuperTableObj
*
pTable
,
const
char
*
tagName
);
int32_t
mgmtSetSchemaFromSuperTable
(
SSchema
*
pSchema
,
SSuperTableObj
*
pTable
);
int32_t
mgmtGetTagsLength
(
SSuperTableObj
*
pSuperTable
,
int32_t
col
);
#ifdef __cplusplus
#ifdef __cplusplus
}
}
...
...
src/mnode/inc/mgmtTable.h
浏览文件 @
03da5549
...
@@ -28,20 +28,18 @@ extern "C" {
...
@@ -28,20 +28,18 @@ extern "C" {
int32_t
mgmtInitTables
();
int32_t
mgmtInitTables
();
STableInfo
*
mgmtGetTable
(
char
*
tableId
);
STableInfo
*
mgmtGetTable
(
char
*
tableId
);
STableInfo
*
mgmtGetTableByPos
(
uint32_t
dnodeIp
,
int32_t
vnode
,
int32_t
sid
);
STableInfo
*
mgmtGetTableByPos
(
uint32_t
dnodeIp
,
int32_t
vnode
,
int32_t
sid
);
int32_t
mgmtGetTableMeta
(
SDbObj
*
pDb
,
STableInfo
*
pTable
,
SMeterMeta
*
pMeta
,
bool
usePublicIp
);
int32_t
mgmtRetrieveMetricMeta
(
void
*
pConn
,
char
**
pStart
,
SSuperTableMetaMsg
*
pInfo
);
int32_t
mgmtRetrieveMetricMeta
(
void
*
pConn
,
char
**
pStart
,
SSuperTableMetaMsg
*
pInfo
);
int32_t
mgmtCreateTable
(
SDbObj
*
pDb
,
SCreateTableMsg
*
pCreate
);
int32_t
mgmtCreateTable
(
SDbObj
*
pDb
,
SCreateTableMsg
*
pCreate
);
int32_t
mgmtDropTable
(
SDbObj
*
pDb
,
char
*
meterId
,
int32_t
ignore
);
int32_t
mgmtDropTable
(
SDbObj
*
pDb
,
char
*
meterId
,
int32_t
ignore
);
int32_t
mgmtAlterTable
(
SDbObj
*
pDb
,
SAlterTableMsg
*
pAlter
);
int32_t
mgmtAlterTable
(
SDbObj
*
pDb
,
SAlterTableMsg
*
pAlter
);
int32_t
mgmtGetTableMeta
(
SMeterMeta
*
pMeta
,
SShowObj
*
pShow
,
void
*
pConn
);
int32_t
mgmtGet
Show
TableMeta
(
SMeterMeta
*
pMeta
,
SShowObj
*
pShow
,
void
*
pConn
);
int32_t
mgmtRetrieveTables
(
SShowObj
*
pShow
,
char
*
data
,
int32_t
rows
,
void
*
pConn
);
int32_t
mgmtRetrieve
Show
Tables
(
SShowObj
*
pShow
,
char
*
data
,
int32_t
rows
,
void
*
pConn
);
void
mgmtCleanUpMeters
();
void
mgmtCleanUpMeters
();
void
mgmtAddTableIntoSuperTable
(
SSuperTableObj
*
pStable
);
void
mgmtAddTableIntoSuperTable
(
SSuperTableObj
*
pStable
);
void
mgmtRemoveTableFromSuperTable
(
SSuperTableObj
*
pStable
);
void
mgmtRemoveTableFromSuperTable
(
SSuperTableObj
*
pStable
);
int32_t
mgmtGetSuperTableMeta
(
SMeterMeta
*
pMeta
,
SShowObj
*
pShow
,
void
*
pConn
);
int32_t
mgmtRetrieveSuperTables
(
SShowObj
*
pShow
,
char
*
data
,
int32_t
rows
,
void
*
pConn
);
#ifdef __cplusplus
#ifdef __cplusplus
...
...
src/mnode/src/mgmtChildTable.c
浏览文件 @
03da5549
...
@@ -340,8 +340,8 @@ int32_t mgmtDropChildTable(SDbObj *pDb, SChildTableObj *pTable) {
...
@@ -340,8 +340,8 @@ int32_t mgmtDropChildTable(SDbObj *pDb, SChildTableObj *pTable) {
return
0
;
return
0
;
}
}
SChildTableObj
*
mgmtGetChildTable
(
char
*
tableId
)
{
void
*
mgmtGetChildTable
(
char
*
tableId
)
{
return
(
SChildTableObj
*
)
sdbGetRow
(
tsChildTableSdb
,
tableId
);
return
sdbGetRow
(
tsChildTableSdb
,
tableId
);
}
}
int32_t
mgmtModifyChildTableTagValueByName
(
SChildTableObj
*
pTable
,
char
*
tagName
,
char
*
nContent
)
{
int32_t
mgmtModifyChildTableTagValueByName
(
SChildTableObj
*
pTable
,
char
*
tagName
,
char
*
nContent
)
{
...
@@ -392,3 +392,30 @@ int32_t mgmtModifyChildTableTagValueByName(SChildTableObj *pTable, char *tagName
...
@@ -392,3 +392,30 @@ int32_t mgmtModifyChildTableTagValueByName(SChildTableObj *pTable, char *tagName
return
0
;
return
0
;
}
}
int32_t
mgmtGetChildTableMeta
(
SDbObj
*
pDb
,
SChildTableObj
*
pTable
,
SMeterMeta
*
pMeta
,
bool
usePublicIp
)
{
pMeta
->
uid
=
htobe64
(
pTable
->
uid
);
pMeta
->
sid
=
htonl
(
pTable
->
sid
);
pMeta
->
vgid
=
htonl
(
pTable
->
vgId
);
pMeta
->
sversion
=
htons
(
pTable
->
superTable
->
sversion
);
pMeta
->
precision
=
pDb
->
cfg
.
precision
;
pMeta
->
numOfTags
=
pTable
->
superTable
->
numOfTags
;
pMeta
->
numOfColumns
=
htons
(
pTable
->
superTable
->
numOfColumns
);
pMeta
->
tableType
=
pTable
->
type
;
pMeta
->
contLen
=
sizeof
(
SMeterMeta
)
+
mgmtSetSchemaFromSuperTable
(
pMeta
->
schema
,
pTable
->
superTable
);
SVgObj
*
pVgroup
=
mgmtGetVgroup
(
pTable
->
vgId
);
if
(
pVgroup
==
NULL
)
{
return
TSDB_CODE_INVALID_TABLE
;
}
for
(
int32_t
i
=
0
;
i
<
TSDB_VNODES_SUPPORT
;
++
i
)
{
if
(
usePublicIp
)
{
pMeta
->
vpeerDesc
[
i
].
ip
=
pVgroup
->
vnodeGid
[
i
].
publicIp
;
pMeta
->
vpeerDesc
[
i
].
vnode
=
htonl
(
pVgroup
->
vnodeGid
[
i
].
vnode
);
}
else
{
pMeta
->
vpeerDesc
[
i
].
ip
=
pVgroup
->
vnodeGid
[
i
].
ip
;
pMeta
->
vpeerDesc
[
i
].
vnode
=
htonl
(
pVgroup
->
vnodeGid
[
i
].
vnode
);
}
}
return
TSDB_CODE_SUCCESS
;
}
src/mnode/src/mgmtNormalTable.c
浏览文件 @
03da5549
...
@@ -357,8 +357,8 @@ int32_t mgmtDropNormalTable(SDbObj *pDb, SNormalTableObj *pTable) {
...
@@ -357,8 +357,8 @@ int32_t mgmtDropNormalTable(SDbObj *pDb, SNormalTableObj *pTable) {
return
0
;
return
0
;
}
}
SNormalTableObj
*
mgmtGetNormalTable
(
char
*
tableId
)
{
void
*
mgmtGetNormalTable
(
char
*
tableId
)
{
return
(
SNormalTableObj
*
)
sdbGetRow
(
tsNormalTableSdb
,
tableId
);
return
sdbGetRow
(
tsNormalTableSdb
,
tableId
);
}
}
static
int32_t
mgmtFindNormalTableColumnIndex
(
SNormalTableObj
*
pTable
,
char
*
colName
)
{
static
int32_t
mgmtFindNormalTableColumnIndex
(
SNormalTableObj
*
pTable
,
char
*
colName
)
{
...
@@ -442,3 +442,45 @@ int32_t mgmtDropNormalTableColumnByName(SNormalTableObj *pTable, char *colName)
...
@@ -442,3 +442,45 @@ int32_t mgmtDropNormalTableColumnByName(SNormalTableObj *pTable, char *colName)
return
TSDB_CODE_SUCCESS
;
return
TSDB_CODE_SUCCESS
;
}
}
static
int32_t
mgmtSetSchemaFromNormalTable
(
SSchema
*
pSchema
,
SNormalTableObj
*
pTable
)
{
int32_t
numOfCols
=
pTable
->
numOfColumns
;
for
(
int32_t
i
=
0
;
i
<
numOfCols
;
++
i
)
{
strcpy
(
pSchema
->
name
,
pTable
->
schema
[
i
].
name
);
pSchema
->
type
=
pTable
->
schema
[
i
].
type
;
pSchema
->
bytes
=
htons
(
pTable
->
schema
[
i
].
bytes
);
pSchema
->
colId
=
htons
(
pTable
->
schema
[
i
].
colId
);
pSchema
++
;
}
return
numOfCols
*
sizeof
(
SSchema
);
}
int32_t
mgmtGetNormalTableMeta
(
SDbObj
*
pDb
,
SNormalTableObj
*
pTable
,
SMeterMeta
*
pMeta
,
bool
usePublicIp
)
{
pMeta
->
uid
=
htobe64
(
pTable
->
uid
);
pMeta
->
sid
=
htonl
(
pTable
->
sid
);
pMeta
->
vgid
=
htonl
(
pTable
->
vgId
);
pMeta
->
sversion
=
htons
(
pTable
->
sversion
);
pMeta
->
precision
=
pDb
->
cfg
.
precision
;
pMeta
->
numOfTags
=
0
;
pMeta
->
numOfColumns
=
htons
(
pTable
->
numOfColumns
);
pMeta
->
tableType
=
pTable
->
type
;
pMeta
->
contLen
=
sizeof
(
SMeterMeta
)
+
mgmtSetSchemaFromNormalTable
(
pMeta
->
schema
,
pTable
);
SVgObj
*
pVgroup
=
mgmtGetVgroup
(
pTable
->
vgId
);
if
(
pVgroup
==
NULL
)
{
return
TSDB_CODE_INVALID_TABLE
;
}
for
(
int32_t
i
=
0
;
i
<
TSDB_VNODES_SUPPORT
;
++
i
)
{
if
(
usePublicIp
)
{
pMeta
->
vpeerDesc
[
i
].
ip
=
pVgroup
->
vnodeGid
[
i
].
publicIp
;
pMeta
->
vpeerDesc
[
i
].
vnode
=
htonl
(
pVgroup
->
vnodeGid
[
i
].
vnode
);
}
else
{
pMeta
->
vpeerDesc
[
i
].
ip
=
pVgroup
->
vnodeGid
[
i
].
ip
;
pMeta
->
vpeerDesc
[
i
].
vnode
=
htonl
(
pVgroup
->
vnodeGid
[
i
].
vnode
);
}
}
return
TSDB_CODE_SUCCESS
;
}
src/mnode/src/mgmtShell.c
浏览文件 @
03da5549
...
@@ -24,17 +24,22 @@
...
@@ -24,17 +24,22 @@
#include "mnode.h"
#include "mnode.h"
#include "mgmtAcct.h"
#include "mgmtAcct.h"
#include "mgmtBalance.h"
#include "mgmtBalance.h"
#include "mgmtChildTable.h"
#include "mgmtConn.h"
#include "mgmtConn.h"
#include "mgmtDb.h"
#include "mgmtDb.h"
#include "mgmtDnode.h"
#include "mgmtDnode.h"
#include "mgmtGrant.h"
#include "mgmtGrant.h"
#include "mgmtMnode.h"
#include "mgmtMnode.h"
#include "mgmtNormalTable.h"
#include "mgmtProfile.h"
#include "mgmtProfile.h"
#include "mgmtShell.h"
#include "mgmtShell.h"
#include "mgmtStreamTable.h"
#include "mgmtSuperTable.h"
#include "mgmtTable.h"
#include "mgmtTable.h"
#include "mgmtUser.h"
#include "mgmtUser.h"
#include "mgmtVgroup.h"
#include "mgmtVgroup.h"
#define MAX_LEN_OF_METER_META (sizeof(SMultiMeterMeta) + sizeof(SSchema) * TSDB_MAX_COLUMNS + sizeof(SSchema) * TSDB_MAX_TAGS + TSDB_MAX_TAGS_LEN)
#define MAX_LEN_OF_METER_META (sizeof(SMultiMeterMeta) + sizeof(SSchema) * TSDB_MAX_COLUMNS + sizeof(SSchema) * TSDB_MAX_TAGS + TSDB_MAX_TAGS_LEN)
typedef
int32_t
(
*
GetMateFp
)(
SMeterMeta
*
pMeta
,
SShowObj
*
pShow
,
void
*
pConn
);
typedef
int32_t
(
*
GetMateFp
)(
SMeterMeta
*
pMeta
,
SShowObj
*
pShow
,
void
*
pConn
);
...
@@ -113,407 +118,135 @@ void mgmtCleanUpShell() {
...
@@ -113,407 +118,135 @@ void mgmtCleanUpShell() {
}
}
}
}
static
void
mgmtSetSchemaFromMeters
(
SSchema
*
pSchema
,
STabObj
*
pMeterObj
,
uint32_t
numOfCols
)
{
int32_t
mgmtProcessMeterMetaMsg
(
void
*
pCont
,
int32_t
contLen
,
void
*
ahandle
)
{
SSchema
*
pMeterSchema
=
(
SSchema
*
)(
pMeterObj
->
schema
);
SRpcConnInfo
connInfo
;
for
(
int32_t
i
=
0
;
i
<
numOfCols
;
++
i
)
{
rpcGetConnInfo
(
ahandle
,
&
connInfo
);
strcpy
(
pSchema
->
name
,
pMeterSchema
[
i
].
name
);
pSchema
->
type
=
pMeterSchema
[
i
].
type
;
bool
usePublicIp
=
(
connInfo
.
serverIp
==
tsPublicIpInt
);
pSchema
->
bytes
=
htons
(
pMeterSchema
[
i
].
bytes
);
SUserObj
*
pUser
=
mgmtGetUser
(
connInfo
.
user
);
pSchema
->
colId
=
htons
(
pMeterSchema
[
i
].
colId
);
if
(
pUser
==
NULL
)
{
pSchema
++
;
rpcSendResponse
(
ahandle
,
TSDB_CODE_INVALID_USER
,
NULL
,
0
);
return
TSDB_CODE_INVALID_USER
;
}
STableInfoMsg
*
pInfo
=
pCont
;
pInfo
->
createFlag
=
htons
(
pInfo
->
createFlag
);
SDbObj
*
pDb
=
mgmtGetDb
(
pInfo
->
db
);
if
(
pDb
==
NULL
||
pDb
->
dropStatus
!=
TSDB_DB_STATUS_READY
)
{
rpcSendResponse
(
ahandle
,
TSDB_CODE_INVALID_DB
,
NULL
,
0
);
return
TSDB_CODE_INVALID_DB
;
}
}
}
static
uint32_t
mgmtSetMeterTagValue
(
char
*
pTags
,
STabObj
*
pMetric
,
STabObj
*
pMeterObj
)
{
STableInfo
*
pTable
=
mgmtGetTable
(
pInfo
->
tableId
);
SSchema
*
pTagSchema
=
(
SSchema
*
)(
pMetric
->
schema
+
pMetric
->
numOfColumns
*
sizeof
(
SSchema
));
char
*
tagVal
=
pMeterObj
->
pTagData
+
TSDB_TABLE_ID_LEN
;
// tag start position
// on demand create table from super table if meter does not exists
if
(
pTable
==
NULL
&&
pInfo
->
createFlag
==
1
)
{
// write operation needs to redirect to master mnode
if
(
mgmtCheckRedirectMsg
(
ahandle
)
!=
0
)
{
return
TSDB_CODE_REDIRECT
;
}
uint32_t
tagsLen
=
0
;
SCreateTableMsg
*
pCreateMsg
=
calloc
(
1
,
sizeof
(
SCreateTableMsg
)
+
sizeof
(
STagData
));
for
(
int32_t
i
=
0
;
i
<
pMetric
->
numOfTags
;
++
i
)
{
if
(
pCreateMsg
==
NULL
)
{
tagsLen
+=
pTagSchema
[
i
].
bytes
;
rpcSendResponse
(
ahandle
,
TSDB_CODE_SERV_OUT_OF_MEMORY
,
NULL
,
0
);
return
TSDB_CODE_SERV_OUT_OF_MEMORY
;
}
memcpy
(
pCreateMsg
->
schema
,
pInfo
->
tags
,
sizeof
(
STagData
));
strcpy
(
pCreateMsg
->
tableId
,
pInfo
->
tableId
);
int32_t
code
=
mgmtCreateTable
(
pDb
,
pCreateMsg
);
char
stableName
[
TSDB_TABLE_ID_LEN
]
=
{
0
};
strncpy
(
stableName
,
pInfo
->
tags
,
TSDB_TABLE_ID_LEN
);
mTrace
(
"table:%s is auto created by %s from %s, code:%d"
,
pCreateMsg
->
tableId
,
pUser
->
user
,
stableName
,
code
);
tfree
(
pCreateMsg
);
if
(
code
!=
TSDB_CODE_SUCCESS
)
{
rpcSendResponse
(
ahandle
,
code
,
NULL
,
0
);
return
code
;
}
pTable
=
mgmtGetTable
(
pInfo
->
tableId
);
}
}
memcpy
(
pTags
,
tagVal
,
tagsLen
);
if
(
pTable
==
NULL
)
{
return
tagsLen
;
rpcSendResponse
(
ahandle
,
TSDB_CODE_INVALID_TABLE
,
NULL
,
0
);
}
return
TSDB_CODE_INVALID_TABLE
;
}
int32_t
mgmtProcessMeterMetaMsg
(
void
*
pCont
,
int32_t
contLen
,
void
*
ahandle
)
{
SMeterMeta
*
pMeta
=
rpcMallocCont
(
sizeof
(
SMeterMeta
)
+
sizeof
(
SSchema
)
*
TSDB_MAX_COLUMNS
);
// SMeterInfoMsg *pInfo = (SMeterInfoMsg *)pMsg;
int32_t
code
=
mgmtGetTableMeta
(
pDb
,
pTable
,
pMeta
,
usePublicIp
);
// STabObj * pMeterObj = NULL;
// SVgObj * pVgroup = NULL;
if
(
code
==
TSDB_CODE_SUCCESS
)
{
// SMeterMeta * pMeta = NULL;
rpcFreeCont
(
pMeta
);
// SSchema * pSchema = NULL;
rpcSendResponse
(
ahandle
,
TSDB_CODE_SUCCESS
,
NULL
,
0
);
// STaosRsp * pRsp = NULL;
}
else
{
// char * pStart = NULL;
pMeta
->
contLen
=
htons
(
pMeta
->
contLen
);
//
rpcSendResponse
(
ahandle
,
TSDB_CODE_SUCCESS
,
pMeta
,
pMeta
->
contLen
);
// pInfo->createFlag = htons(pInfo->createFlag);
}
//
// int32_t size = sizeof(STaosHeader) + sizeof(STaosRsp) + sizeof(SMeterMeta) + sizeof(SSchema) * TSDB_MAX_COLUMNS +
return
TSDB_CODE_SUCCESS
;
// sizeof(SSchema) * TSDB_MAX_TAGS + TSDB_MAX_TAGS_LEN + TSDB_EXTRA_PAYLOAD_SIZE;
//
// SDbObj *pDb = NULL;
// if (pConn->pDb != NULL) pDb = mgmtGetDb(pConn->pDb->name);
//
// // todo db check should be extracted
// if (pDb == NULL || (pDb != NULL && pDb->dropStatus != TSDB_DB_STATUS_READY)) {
//
// if ((pStart = mgmtAllocMsg(pConn, size, &pMsg, &pRsp)) == NULL) {
// taosSendSimpleRsp(pConn->thandle, TSDB_MSG_TYPE_TABLE_META_RSP, TSDB_CODE_SERV_OUT_OF_MEMORY);
// return 0;
// }
//
// pRsp->code = TSDB_CODE_INVALID_DB;
// pMsg++;
//
// goto _exit_code;
// }
//
// pMeterObj = mgmtGetTable(pInfo->meterId);
//
// // on demand create table from super table if meter does not exists
// if (pMeterObj == NULL && pInfo->createFlag == 1) {
// // write operation needs to redirect to master mnode
// if (mgmtCheckRedirectMsg(pConn, TSDB_MSG_TYPE_TABLE_META_RSP) != 0) {
// return 0;
// }
//
// SCreateTableMsg *pCreateMsg = calloc(1, sizeof(SCreateTableMsg) + sizeof(STagData));
// if (pCreateMsg == NULL) {
// taosSendSimpleRsp(pConn->thandle, TSDB_MSG_TYPE_TABLE_META_RSP, TSDB_CODE_SERV_OUT_OF_MEMORY);
// return 0;
// }
//
// memcpy(pCreateMsg->schema, pInfo->tags, sizeof(STagData));
// strcpy(pCreateMsg->meterId, pInfo->meterId);
//
// SDbObj* pMeterDb = mgmtGetDbByTableId(pCreateMsg->meterId);
// mTrace("table:%s, pConnDb:%p, pConnDbName:%s, pMeterDb:%p, pMeterDbName:%s",
// pCreateMsg->meterId, pDb, pDb->name, pMeterDb, pMeterDb->name);
// assert(pDb == pMeterDb);
//
// int32_t code = mgmtCreateTable(pDb, pCreateMsg);
//
// char stableName[TSDB_TABLE_ID_LEN] = {0};
// strncpy(stableName, pInfo->tags, TSDB_TABLE_ID_LEN);
// mTrace("table:%s is automatically created by %s from %s, code:%d", pCreateMsg->meterId, pConn->pUser->user,
// stableName, code);
//
// tfree(pCreateMsg);
//
// if (code != TSDB_CODE_SUCCESS) {
// if ((pStart = mgmtAllocMsg(pConn, size, &pMsg, &pRsp)) == NULL) {
// taosSendSimpleRsp(pConn->thandle, TSDB_MSG_TYPE_TABLE_META_RSP, TSDB_CODE_SERV_OUT_OF_MEMORY);
// return 0;
// }
//
// pRsp->code = code;
// pMsg++;
//
// goto _exit_code;
// }
//
// pMeterObj = mgmtGetTable(pInfo->meterId);
// }
//
// if ((pStart = mgmtAllocMsg(pConn, size, &pMsg, &pRsp)) == NULL) {
// taosSendSimpleRsp(pConn->thandle, TSDB_MSG_TYPE_TABLE_META_RSP, TSDB_CODE_SERV_OUT_OF_MEMORY);
// return 0;
// }
//
// if (pMeterObj == NULL) {
// if (pDb)
// pRsp->code = TSDB_CODE_INVALID_TABLE;
// else
// pRsp->code = TSDB_CODE_DB_NOT_SELECTED;
// pMsg++;
// } else {
// mTrace("%s, uid:%" PRIu64 " meter meta is retrieved", pInfo->meterId, pMeterObj->uid);
// pRsp->code = 0;
// pMsg += sizeof(STaosRsp);
// *pMsg = TSDB_IE_TYPE_META;
// pMsg++;
//
// pMeta = (SMeterMeta *)pMsg;
// pMeta->uid = htobe64(pMeterObj->uid);
// pMeta->sid = htonl(pMeterObj->gid.sid);
// pMeta->vgid = htonl(pMeterObj->gid.vgId);
// pMeta->sversion = htons(pMeterObj->sversion);
//
// pMeta->precision = pDb->cfg.precision;
//
// pMeta->numOfTags = pMeterObj->numOfTags;
// pMeta->numOfColumns = htons(pMeterObj->numOfColumns);
// pMeta->tableType = pMeterObj->tableType;
//
// pMsg += sizeof(SMeterMeta);
// pSchema = (SSchema *)pMsg; // schema locates at the end of SMeterMeta struct
//
// if (mgmtTableCreateFromSuperTable(pMeterObj)) {
// assert(pMeterObj->numOfTags == 0);
//
// STabObj *pMetric = mgmtGetTable(pMeterObj->pTagData);
// uint32_t numOfTotalCols = (uint32_t)pMetric->numOfTags + pMetric->numOfColumns;
//
// pMeta->numOfTags = pMetric->numOfTags; // update the numOfTags info
// mgmtSetSchemaFromMeters(pSchema, pMetric, numOfTotalCols);
// pMsg += numOfTotalCols * sizeof(SSchema);
//
// // for meters created from metric, we need the metric tag schema to parse the tag data
// int32_t tagsLen = mgmtSetMeterTagValue(pMsg, pMetric, pMeterObj);
// pMsg += tagsLen;
// } else {
// /*
// * for metrics, or meters that are not created from metric, set the schema directly
// * for meters created from metric, we use the schema of metric instead
// */
// uint32_t numOfTotalCols = (uint32_t)pMeterObj->numOfTags + pMeterObj->numOfColumns;
// mgmtSetSchemaFromMeters(pSchema, pMeterObj, numOfTotalCols);
// pMsg += numOfTotalCols * sizeof(SSchema);
// }
//
// if (mgmtIsNormalTable(pMeterObj)) {
// pVgroup = mgmtGetVgroup(pMeterObj->gid.vgId);
// if (pVgroup == NULL) {
// pRsp->code = TSDB_CODE_INVALID_TABLE;
// goto _exit_code;
// }
// for (int32_t i = 0; i < TSDB_VNODES_SUPPORT; ++i) {
// if (pConn->usePublicIp) {
// pMeta->vpeerDesc[i].ip = pVgroup->vnodeGid[i].publicIp;
// pMeta->vpeerDesc[i].vnode = htonl(pVgroup->vnodeGid[i].vnode);
// } else {
// pMeta->vpeerDesc[i].ip = pVgroup->vnodeGid[i].ip;
// pMeta->vpeerDesc[i].vnode = htonl(pVgroup->vnodeGid[i].vnode);
// }
// }
// }
// }
//
//_exit_code:
// msgLen = pMsg - pStart;
//
// taosSendMsgToPeer(pConn->thandle, pStart, msgLen);
//
// return msgLen;
return
0
;
}
}
/**
* multi meter meta rsp pkg format:
* | STaosRsp | ieType | SMultiMeterInfoMsg | SMeterMeta0 | SSchema0 | SMeterMeta1 | SSchema1 | SMeterMeta2 | SSchema2
* 1B 1B 4B
*
* | STaosHeader | STaosRsp | ieType | SMultiMeterInfoMsg | SMeterMeta0 | SSchema0 | SMeterMeta1 | SSchema1 | ......................|
* ^ ^ ^
* |<--------------------------------------size-----------------------------------------------|---------------------->|
* | | |
* pStart pCurMeter pTail
**/
int32_t
mgmtProcessMultiMeterMetaMsg
(
void
*
pCont
,
int32_t
contLen
,
void
*
ahandle
)
{
int32_t
mgmtProcessMultiMeterMetaMsg
(
void
*
pCont
,
int32_t
contLen
,
void
*
ahandle
)
{
// SDbObj * pDbObj = NULL;
SRpcConnInfo
connInfo
;
// STabObj * pMeterObj = NULL;
rpcGetConnInfo
(
ahandle
,
&
connInfo
);
// SVgObj * pVgroup = NULL;
// SMultiMeterMeta * pMeta = NULL;
// SSchema * pSchema = NULL;
// STaosRsp * pRsp = NULL;
// char * pStart = NULL;
//
// SMultiMeterInfoMsg * pInfo = (SMultiMeterInfoMsg *)pMsg;
// char * str = pMsg + sizeof(SMultiMeterInfoMsg);
// pInfo->numOfMeters = htonl(pInfo->numOfMeters);
//
// int32_t size = 4*1024*1024; // first malloc 4 MB, subsequent reallocation as twice
//
// char *pNewMsg;
// if ((pStart = mgmtForMultiAllocMsg(pConn, size, &pNewMsg, &pRsp)) == NULL) {
// taosSendSimpleRsp(pConn->thandle, TSDB_MSG_TYPE_MULTI_TABLE_META_RSP, TSDB_CODE_SERV_OUT_OF_MEMORY);
// return 0;
// }
//
// int32_t totalNum = 0;
// char tblName[TSDB_TABLE_ID_LEN];
// char* nextStr;
//
// char* pCurMeter = pStart + sizeof(STaosRsp) + sizeof(SMultiMeterInfoMsg) + 1; // 1: ie type byte
// char* pTail = pStart + size;
//
// while (str - pMsg < msgLen) {
// nextStr = strchr(str, ',');
// if (nextStr == NULL) {
// break;
// }
//
// memcpy(tblName, str, nextStr - str);
// tblName[nextStr - str] = '\0';
// str = nextStr + 1;
//
// // judge whether the remaining memory is adequate
// if ((pTail - pCurMeter) < MAX_LEN_OF_METER_META) {
// char* pMsgHdr = pStart - sizeof(STaosHeader);
// size *= 2;
// pMsgHdr = (char*)realloc(pMsgHdr, size);
// if (NULL == pMsgHdr) {
// char* pTmp = pStart - sizeof(STaosHeader);
// tfree(pTmp);
// taosSendSimpleRsp(pConn->thandle, TSDB_MSG_TYPE_MULTI_TABLE_META_RSP, TSDB_CODE_SERV_OUT_OF_MEMORY);
// break;
// }
//
// pCurMeter = (char*)pMsgHdr + sizeof(STaosHeader) + (pCurMeter - pStart);
// pStart = (char*)pMsgHdr + sizeof(STaosHeader);
// pNewMsg = pStart;
// pRsp = (STaosRsp *)pStart;
// pTail = pMsgHdr + size;
// }
//
// // get meter schema, and fill into resp payload
// pMeterObj = mgmtGetTable(tblName);
// pDbObj = mgmtGetDbByTableId(tblName);
//
// if (pMeterObj == NULL || (pDbObj == NULL)) {
// continue;
// } else {
// mTrace("%s, uid:%" PRIu64 " sversion:%d meter meta is retrieved", tblName, pMeterObj->uid, pMeterObj->sversion);
// pMeta = (SMultiMeterMeta *)pCurMeter;
//
// memcpy(pMeta->meterId, tblName, strlen(tblName));
// pMeta->meta.uid = htobe64(pMeterObj->uid);
// pMeta->meta.sid = htonl(pMeterObj->gid.sid);
// pMeta->meta.vgid = htonl(pMeterObj->gid.vgId);
// pMeta->meta.sversion = htons(pMeterObj->sversion);
// pMeta->meta.precision = pDbObj->cfg.precision;
// pMeta->meta.numOfTags = pMeterObj->numOfTags;
// pMeta->meta.numOfColumns = htons(pMeterObj->numOfColumns);
// pMeta->meta.tableType = pMeterObj->tableType;
//
// pCurMeter += sizeof(SMultiMeterMeta);
// pSchema = (SSchema *)pCurMeter; // schema locates at the end of SMeterMeta struct
//
// if (mgmtTableCreateFromSuperTable(pMeterObj)) {
// assert(pMeterObj->numOfTags == 0);
//
// STabObj *pMetric = mgmtGetTable(pMeterObj->pTagData);
// uint32_t numOfTotalCols = (uint32_t)pMetric->numOfTags + pMetric->numOfColumns;
//
// pMeta->meta.numOfTags = pMetric->numOfTags; // update the numOfTags info
// mgmtSetSchemaFromMeters(pSchema, pMetric, numOfTotalCols);
// pCurMeter += numOfTotalCols * sizeof(SSchema);
//
// // for meters created from metric, we need the metric tag schema to parse the tag data
// int32_t tagsLen = mgmtSetMeterTagValue(pCurMeter, pMetric, pMeterObj);
// pCurMeter += tagsLen;
// } else {
// /*
// * for metrics, or meters that are not created from metric, set the schema directly
// * for meters created from metric, we use the schema of metric instead
// */
// uint32_t numOfTotalCols = (uint32_t)pMeterObj->numOfTags + pMeterObj->numOfColumns;
// mgmtSetSchemaFromMeters(pSchema, pMeterObj, numOfTotalCols);
// pCurMeter += numOfTotalCols * sizeof(SSchema);
// }
//
// if (mgmtIsNormalTable(pMeterObj)) {
// pVgroup = mgmtGetVgroup(pMeterObj->gid.vgId);
// if (pVgroup == NULL) {
// pRsp->code = TSDB_CODE_INVALID_TABLE;
// pNewMsg++;
// mError("%s, uid:%" PRIu64 " sversion:%d vgId:%d pVgroup is NULL", tblName, pMeterObj->uid, pMeterObj->sversion,
// pMeterObj->gid.vgId);
// goto _error_exit_code;
// }
//
// for (int32_t i = 0; i < TSDB_VNODES_SUPPORT; ++i) {
// if (pConn->usePublicIp) {
// pMeta->meta.vpeerDesc[i].ip = pVgroup->vnodeGid[i].publicIp;
// pMeta->meta.vpeerDesc[i].vnode = htonl(pVgroup->vnodeGid[i].vnode);
// } else {
// pMeta->meta.vpeerDesc[i].ip = pVgroup->vnodeGid[i].ip;
// pMeta->meta.vpeerDesc[i].vnode = htonl(pVgroup->vnodeGid[i].vnode);
// }
// }
// }
// }
//
// totalNum++;
// if (totalNum > pInfo->numOfMeters) {
// pNewMsg++;
// break;
// }
// }
//
// // fill rsp code, ieType
// msgLen = pCurMeter - pNewMsg;
//
// pRsp->code = 0;
// pNewMsg += sizeof(STaosRsp);
// *pNewMsg = TSDB_IE_TYPE_META;
// pNewMsg++;
//
// SMultiMeterInfoMsg *pRspInfo = (SMultiMeterInfoMsg *)pNewMsg;
//
// pRspInfo->numOfMeters = htonl(totalNum);
// goto _exit_code;
//
//_error_exit_code:
// msgLen = pNewMsg - pStart;
//
//_exit_code:
// taosSendMsgToPeer(pConn->thandle, pStart, msgLen);
//
// return msgLen;
return
0
;
}
int32_t
mgmtProcessMetricMetaMsg
(
void
*
pCont
,
int32_t
contLen
,
void
*
ahandle
)
{
bool
usePublicIp
=
(
connInfo
.
serverIp
==
tsPublicIpInt
);
// SSuperTableMetaMsg *pSuperTableMetaMsg = (SSuperTableMetaMsg *)pMsg;
SUserObj
*
pUser
=
mgmtGetUser
(
connInfo
.
user
);
// STabObj * pMetric;
if
(
pUser
==
NULL
)
{
// STaosRsp * pRsp;
rpcSendResponse
(
ahandle
,
TSDB_CODE_INVALID_USER
,
NULL
,
0
);
// char * pStart;
return
TSDB_CODE_INVALID_USER
;
//
}
// pSuperTableMetaMsg->numOfMeters = htonl(pSuperTableMetaMsg->numOfMeters);
//
SMultiTableInfoMsg
*
pInfo
=
pCont
;
// pSuperTableMetaMsg->join = htonl(pSuperTableMetaMsg->join);
pInfo
->
numOfTables
=
htonl
(
pInfo
->
numOfTables
);
// pSuperTableMetaMsg->joinCondLen = htonl(pSuperTableMetaMsg->joinCondLen);
//
int32_t
totalMallocLen
=
4
*
1024
*
1024
;
// first malloc 4 MB, subsequent reallocation as twice
// for (int32_t i = 0; i < pSuperTableMetaMsg->numOfMeters; ++i) {
SMultiMeterMeta
*
pMultiMeta
=
rpcMallocCont
(
totalMallocLen
);
// pSuperTableMetaMsg->metaElem[i] = htonl(pSuperTableMetaMsg->metaElem[i]);
if
(
pMultiMeta
==
NULL
)
{
// }
rpcSendResponse
(
ahandle
,
TSDB_CODE_SERV_OUT_OF_MEMORY
,
NULL
,
0
);
//
return
TSDB_CODE_SERV_OUT_OF_MEMORY
;
// SMetricMetaElemMsg *pElem = (SMetricMetaElemMsg *)(((char *)pSuperTableMetaMsg) + pSuperTableMetaMsg->metaElem[0]);
}
// pMetric = mgmtGetTable(pElem->meterId);
//
pMultiMeta
->
contLen
=
sizeof
(
SMultiMeterMeta
);
// SDbObj *pDb = NULL;
pMultiMeta
->
numOfTables
=
0
;
// if (pConn->pDb != NULL) pDb = mgmtGetDb(pConn->pDb->name);
//
for
(
int
t
=
0
;
t
<
pInfo
->
numOfTables
;
++
t
)
{
// if (pMetric == NULL || (pDb != NULL && pDb->dropStatus != TSDB_DB_STATUS_READY)) {
char
*
tableId
=
(
char
*
)(
pInfo
->
tableIds
+
t
*
TSDB_TABLE_ID_LEN
);
// pStart = taosBuildRspMsg(pConn->thandle, TSDB_MSG_TYPE_STABLE_META_RSP);
STableInfo
*
pTable
=
mgmtGetTable
(
tableId
);
// if (pStart == NULL) {
if
(
pTable
==
NULL
)
continue
;
// taosSendSimpleRsp(pConn->thandle, TSDB_MSG_TYPE_STABLE_META_RSP, TSDB_CODE_SERV_OUT_OF_MEMORY);
// return 0;
SDbObj
*
pDb
=
mgmtGetDbByTableId
(
tableId
);
// }
if
(
pDb
==
NULL
)
continue
;
//
// pMsg = pStart;
int
availLen
=
totalMallocLen
-
pMultiMeta
->
contLen
;
// pRsp = (STaosRsp *)pMsg;
if
(
availLen
<=
sizeof
(
SMeterMeta
)
+
sizeof
(
SSchema
)
*
TSDB_MAX_COLUMNS
)
{
// if (pDb)
//TODO realloc
// pRsp->code = TSDB_CODE_INVALID_TABLE;
//totalMallocLen *= 2;
// else
//pMultiMeta = rpcReMalloc(pMultiMeta, totalMallocLen);
// pRsp->code = TSDB_CODE_DB_NOT_SELECTED;
//if (pMultiMeta == NULL) {
// pMsg++;
/// rpcSendResponse(ahandle, TSDB_CODE_SERV_OUT_OF_MEMORY, NULL, 0);
//
// return TSDB_CODE_SERV_OUT_OF_MEMORY;
// msgLen = pMsg - pStart;
//} else {
// } else {
// t--;
// msgLen = mgmtRetrieveMetricMeta(pConn, &pStart, pSuperTableMetaMsg);
// continue;
// if (msgLen <= 0) {
//}
// taosSendSimpleRsp(pConn->thandle, TSDB_MSG_TYPE_STABLE_META_RSP, TSDB_CODE_SERV_OUT_OF_MEMORY);
}
// return 0;
// }
SMeterMeta
*
pMeta
=
(
SMeterMeta
*
)(
pMultiMeta
->
metas
+
pMultiMeta
->
contLen
);
// }
int32_t
code
=
mgmtGetTableMeta
(
pDb
,
pTable
,
pMeta
,
usePublicIp
);
//
if
(
code
==
TSDB_CODE_SUCCESS
)
{
// taosSendMsgToPeer(pConn->thandle, pStart, msgLen);
pMultiMeta
->
numOfTables
++
;
//
pMultiMeta
->
contLen
+=
pMeta
->
contLen
;
// return msgLen;
}
return
0
;
}
rpcSendResponse
(
ahandle
,
TSDB_CODE_SUCCESS
,
pMultiMeta
,
pMultiMeta
->
contLen
);
return
TSDB_CODE_SUCCESS
;
}
}
int32_t
mgmtProcessCreateDbMsg
(
void
*
pCont
,
int32_t
contLen
,
void
*
ahandle
)
{
int32_t
mgmtProcessCreateDbMsg
(
void
*
pCont
,
int32_t
contLen
,
void
*
ahandle
)
{
...
@@ -873,40 +606,40 @@ int32_t mgmtProcessDropDbMsg(void *pCont, int32_t contLen, void *ahandle) {
...
@@ -873,40 +606,40 @@ int32_t mgmtProcessDropDbMsg(void *pCont, int32_t contLen, void *ahandle) {
static
void
mgmtInitShowMsgFp
()
{
static
void
mgmtInitShowMsgFp
()
{
mgmtGetMetaFp
=
(
GetMateFp
*
)
malloc
(
TSDB_MGMT_TABLE_MAX
*
sizeof
(
GetMateFp
));
mgmtGetMetaFp
=
(
GetMateFp
*
)
malloc
(
TSDB_MGMT_TABLE_MAX
*
sizeof
(
GetMateFp
));
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_ACCT
]
=
mgmtGetAcctMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_ACCT
]
=
mgmtGetAcctMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_USER
]
=
mgmtGetUserMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_USER
]
=
mgmtGetUserMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_DB
]
=
mgmtGetDbMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_DB
]
=
mgmtGetDbMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_TABLE
]
=
mgmtGet
TableMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_TABLE
]
=
mgmtGetShow
TableMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_DNODE
]
=
mgmtGetDnodeMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_DNODE
]
=
mgmtGetDnodeMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_MNODE
]
=
mgmtGetMnodeMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_MNODE
]
=
mgmtGetMnodeMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_VGROUP
]
=
mgmtGetVgroupMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_VGROUP
]
=
mgmtGetVgroupMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_METRIC
]
=
mgmtGet
SuperTableMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_METRIC
]
=
mgmtGetShow
SuperTableMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_MODULE
]
=
mgmtGetModuleMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_MODULE
]
=
mgmtGetModuleMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_QUERIES
]
=
mgmtGetQueryMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_QUERIES
]
=
mgmtGetQueryMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_STREAMS
]
=
mgmtGetStreamMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_STREAMS
]
=
mgmtGetStreamMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_CONFIGS
]
=
mgmtGetConfigMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_CONFIGS
]
=
mgmtGetConfigMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_CONNS
]
=
mgmtGetConnsMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_CONNS
]
=
mgmtGetConnsMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_SCORES
]
=
mgmtGetScoresMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_SCORES
]
=
mgmtGetScoresMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_GRANTS
]
=
mgmtGetGrantsMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_GRANTS
]
=
mgmtGetGrantsMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_VNODES
]
=
mgmtGetVnodeMeta
;
mgmtGetMetaFp
[
TSDB_MGMT_TABLE_VNODES
]
=
mgmtGetVnodeMeta
;
mgmtRetrieveFp
=
(
RetrieveMetaFp
*
)
malloc
(
TSDB_MGMT_TABLE_MAX
*
sizeof
(
RetrieveMetaFp
));
mgmtRetrieveFp
=
(
RetrieveMetaFp
*
)
malloc
(
TSDB_MGMT_TABLE_MAX
*
sizeof
(
RetrieveMetaFp
));
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_ACCT
]
=
mgmtRetrieveAccts
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_ACCT
]
=
mgmtRetrieveAccts
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_USER
]
=
mgmtRetrieveUsers
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_USER
]
=
mgmtRetrieveUsers
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_DB
]
=
mgmtRetrieveDbs
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_DB
]
=
mgmtRetrieveDbs
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_TABLE
]
=
mgmtRetrieve
Tables
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_TABLE
]
=
mgmtRetrieveShow
Tables
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_DNODE
]
=
mgmtRetrieveDnodes
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_DNODE
]
=
mgmtRetrieveDnodes
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_MNODE
]
=
mgmtRetrieveMnodes
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_MNODE
]
=
mgmtRetrieveMnodes
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_VGROUP
]
=
mgmtRetrieveVgroups
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_VGROUP
]
=
mgmtRetrieveVgroups
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_METRIC
]
=
mgmtRetrieve
SuperTables
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_METRIC
]
=
mgmtRetrieveShow
SuperTables
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_MODULE
]
=
mgmtRetrieveModules
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_MODULE
]
=
mgmtRetrieveModules
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_QUERIES
]
=
mgmtRetrieveQueries
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_QUERIES
]
=
mgmtRetrieveQueries
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_STREAMS
]
=
mgmtRetrieveStreams
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_STREAMS
]
=
mgmtRetrieveStreams
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_CONFIGS
]
=
mgmtRetrieveConfigs
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_CONFIGS
]
=
mgmtRetrieveConfigs
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_CONNS
]
=
mgmtRetrieveConns
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_CONNS
]
=
mgmtRetrieveConns
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_SCORES
]
=
mgmtRetrieveScores
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_SCORES
]
=
mgmtRetrieveScores
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_GRANTS
]
=
mgmtRetrieveGrants
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_GRANTS
]
=
mgmtRetrieveGrants
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_VNODES
]
=
mgmtRetrieveVnodes
;
mgmtRetrieveFp
[
TSDB_MGMT_TABLE_VNODES
]
=
mgmtRetrieveVnodes
;
}
}
int32_t
mgmtProcessShowMsg
(
void
*
pCont
,
int32_t
contLen
,
void
*
ahandle
)
{
int32_t
mgmtProcessShowMsg
(
void
*
pCont
,
int32_t
contLen
,
void
*
ahandle
)
{
...
@@ -1084,9 +817,9 @@ int32_t mgmtProcessDropTableMsg(void *pCont, int32_t contLen, void *ahandle) {
...
@@ -1084,9 +817,9 @@ int32_t mgmtProcessDropTableMsg(void *pCont, int32_t contLen, void *ahandle) {
}
else
{
}
else
{
SDbObj
*
pDb
=
mgmtGetDb
(
pDrop
->
db
);
SDbObj
*
pDb
=
mgmtGetDb
(
pDrop
->
db
);
if
(
pDb
)
{
if
(
pDb
)
{
code
=
mgmtDropTable
(
pDb
,
pDrop
->
meter
Id
,
pDrop
->
igNotExists
);
code
=
mgmtDropTable
(
pDb
,
pDrop
->
table
Id
,
pDrop
->
igNotExists
);
if
(
code
==
TSDB_CODE_SUCCESS
)
{
if
(
code
==
TSDB_CODE_SUCCESS
)
{
mTrace
(
"table:%s is dropped by user:%s"
,
pDrop
->
meter
Id
,
pUser
->
user
);
mTrace
(
"table:%s is dropped by user:%s"
,
pDrop
->
table
Id
,
pUser
->
user
);
}
}
}
else
{
}
else
{
code
=
TSDB_CODE_DB_NOT_SELECTED
;
code
=
TSDB_CODE_DB_NOT_SELECTED
;
...
@@ -1310,14 +1043,14 @@ connect_over:
...
@@ -1310,14 +1043,14 @@ connect_over:
* check if we need to add mgmtProcessMeterMetaMsg into tranQueue, which will be executed one-by-one.
* check if we need to add mgmtProcessMeterMetaMsg into tranQueue, which will be executed one-by-one.
*/
*/
static
bool
mgmtCheckMeterMetaMsgType
(
void
*
pMsg
)
{
static
bool
mgmtCheckMeterMetaMsgType
(
void
*
pMsg
)
{
S
MeterInfoMsg
*
pInfo
=
(
SMeter
InfoMsg
*
)
pMsg
;
S
TableInfoMsg
*
pInfo
=
(
STable
InfoMsg
*
)
pMsg
;
int16_t
autoCreate
=
htons
(
pInfo
->
createFlag
);
int16_t
autoCreate
=
htons
(
pInfo
->
createFlag
);
STableInfo
*
pTable
=
mgmtGetTable
(
pInfo
->
meter
Id
);
STableInfo
*
pTable
=
mgmtGetTable
(
pInfo
->
table
Id
);
// If table does not exists and autoCreate flag is set, we add the handler into task queue
// If table does not exists and autoCreate flag is set, we add the handler into task queue
bool
addIntoTranQueue
=
(
pTable
==
NULL
&&
autoCreate
==
1
);
bool
addIntoTranQueue
=
(
pTable
==
NULL
&&
autoCreate
==
1
);
if
(
addIntoTranQueue
)
{
if
(
addIntoTranQueue
)
{
mTrace
(
"table:%s auto created task added"
,
pInfo
->
meter
Id
);
mTrace
(
"table:%s auto created task added"
,
pInfo
->
table
Id
);
}
}
return
addIntoTranQueue
;
return
addIntoTranQueue
;
...
@@ -1380,10 +1113,10 @@ void mgmtInitProcessShellMsg() {
...
@@ -1380,10 +1113,10 @@ void mgmtInitProcessShellMsg() {
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_KILL_STREAM
]
=
mgmtProcessKillStreamMsg
;
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_KILL_STREAM
]
=
mgmtProcessKillStreamMsg
;
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_KILL_CONNECTION
]
=
mgmtProcessKillConnectionMsg
;
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_KILL_CONNECTION
]
=
mgmtProcessKillConnectionMsg
;
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_SHOW
]
=
mgmtProcessShowMsg
;
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_SHOW
]
=
mgmtProcessShowMsg
;
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_RETRIEVE
]
=
mgmtProcessRetrieveMsg
;
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_RETRIEVE
]
=
mgmtProcessRetrieveMsg
;
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_TABLE_META
]
=
mgmtProcessMeterMetaMsg
;
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_TABLE_META
]
=
mgmtProcessMeterMetaMsg
;
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_STABLE_META
]
=
mgmtProcessMetricMetaMsg
;
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_MULTI_TABLE_META
]
=
mgmtProcessMultiMeterMetaMsg
;
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_MULTI_TABLE_META
]
=
mgmtProcessMultiMeterMetaMsg
;
mgmtProcessShellMsg
[
TSDB_MSG_TYPE_STABLE_META
]
=
mgmtProcessUnSupportMsg
;
}
}
static
int32_t
mgmtCheckRedirectMsgImp
(
void
*
pConn
)
{
static
int32_t
mgmtCheckRedirectMsgImp
(
void
*
pConn
)
{
...
...
src/mnode/src/mgmtStreamTable.c
浏览文件 @
03da5549
...
@@ -379,6 +379,47 @@ int32_t mgmtDropStreamTable(SDbObj *pDb, SStreamTableObj *pTable) {
...
@@ -379,6 +379,47 @@ int32_t mgmtDropStreamTable(SDbObj *pDb, SStreamTableObj *pTable) {
return
0
;
return
0
;
}
}
SStreamTableObj
*
mgmtGetStreamTable
(
char
*
tableId
)
{
void
*
mgmtGetStreamTable
(
char
*
tableId
)
{
return
(
SStreamTableObj
*
)
sdbGetRow
(
tsStreamTableSdb
,
tableId
);
return
sdbGetRow
(
tsStreamTableSdb
,
tableId
);
}
}
\ No newline at end of file
static
int32_t
mgmtSetSchemaFromStreamTable
(
SSchema
*
pSchema
,
SStreamTableObj
*
pTable
)
{
int32_t
numOfCols
=
pTable
->
numOfColumns
;
for
(
int32_t
i
=
0
;
i
<
numOfCols
;
++
i
)
{
strcpy
(
pSchema
->
name
,
pTable
->
schema
[
i
].
name
);
pSchema
->
type
=
pTable
->
schema
[
i
].
type
;
pSchema
->
bytes
=
htons
(
pTable
->
schema
[
i
].
bytes
);
pSchema
->
colId
=
htons
(
pTable
->
schema
[
i
].
colId
);
pSchema
++
;
}
return
numOfCols
*
sizeof
(
SSchema
);
}
int32_t
mgmtGetStreamTableMeta
(
SDbObj
*
pDb
,
SStreamTableObj
*
pTable
,
SMeterMeta
*
pMeta
,
bool
usePublicIp
)
{
pMeta
->
uid
=
htobe64
(
pTable
->
uid
);
pMeta
->
sid
=
htonl
(
pTable
->
sid
);
pMeta
->
vgid
=
htonl
(
pTable
->
vgId
);
pMeta
->
sversion
=
htons
(
pTable
->
sversion
);
pMeta
->
precision
=
pDb
->
cfg
.
precision
;
pMeta
->
numOfTags
=
0
;
pMeta
->
numOfColumns
=
htons
(
pTable
->
numOfColumns
);
pMeta
->
tableType
=
pTable
->
type
;
pMeta
->
contLen
=
sizeof
(
SMeterMeta
)
+
mgmtSetSchemaFromStreamTable
(
pMeta
->
schema
,
pTable
);
SVgObj
*
pVgroup
=
mgmtGetVgroup
(
pTable
->
vgId
);
if
(
pVgroup
==
NULL
)
{
return
TSDB_CODE_INVALID_TABLE
;
}
for
(
int32_t
i
=
0
;
i
<
TSDB_VNODES_SUPPORT
;
++
i
)
{
if
(
usePublicIp
)
{
pMeta
->
vpeerDesc
[
i
].
ip
=
pVgroup
->
vnodeGid
[
i
].
publicIp
;
pMeta
->
vpeerDesc
[
i
].
vnode
=
htonl
(
pVgroup
->
vnodeGid
[
i
].
vnode
);
}
else
{
pMeta
->
vpeerDesc
[
i
].
ip
=
pVgroup
->
vnodeGid
[
i
].
ip
;
pMeta
->
vpeerDesc
[
i
].
vnode
=
htonl
(
pVgroup
->
vnodeGid
[
i
].
vnode
);
}
}
return
TSDB_CODE_SUCCESS
;
}
src/mnode/src/mgmtSuperTable.c
浏览文件 @
03da5549
...
@@ -235,8 +235,8 @@ int32_t mgmtDropSuperTable(SDbObj *pDb, SSuperTableObj *pSuperTable) {
...
@@ -235,8 +235,8 @@ int32_t mgmtDropSuperTable(SDbObj *pDb, SSuperTableObj *pSuperTable) {
return
sdbDeleteRow
(
tsSuperTableSdb
,
pSuperTable
);
return
sdbDeleteRow
(
tsSuperTableSdb
,
pSuperTable
);
}
}
SSuperTableObj
*
mgmtGetSuperTable
(
char
*
tableId
)
{
void
*
mgmtGetSuperTable
(
char
*
tableId
)
{
return
(
SSuperTableObj
*
)
sdbGetRow
(
tsSuperTableSdb
,
tableId
);
return
sdbGetRow
(
tsSuperTableSdb
,
tableId
);
}
}
int32_t
mgmtFindSuperTableTagIndex
(
SSuperTableObj
*
pStable
,
const
char
*
tagName
)
{
int32_t
mgmtFindSuperTableTagIndex
(
SSuperTableObj
*
pStable
,
const
char
*
tagName
)
{
...
@@ -457,7 +457,7 @@ int32_t mgmtDropSuperTableColumnByName(SSuperTableObj *pStable, char *colName) {
...
@@ -457,7 +457,7 @@ int32_t mgmtDropSuperTableColumnByName(SSuperTableObj *pStable, char *colName) {
return
TSDB_CODE_SUCCESS
;
return
TSDB_CODE_SUCCESS
;
}
}
int32_t
mgmtGetSuperTableMeta
(
SMeterMeta
*
pMeta
,
SShowObj
*
pShow
,
void
*
pConn
)
{
int32_t
mgmtGetS
howS
uperTableMeta
(
SMeterMeta
*
pMeta
,
SShowObj
*
pShow
,
void
*
pConn
)
{
// int32_t cols = 0;
// int32_t cols = 0;
//
//
// SDbObj *pDb = NULL;
// SDbObj *pDb = NULL;
...
@@ -509,7 +509,7 @@ int32_t mgmtGetSuperTableMeta(SMeterMeta *pMeta, SShowObj *pShow, void *pConn) {
...
@@ -509,7 +509,7 @@ int32_t mgmtGetSuperTableMeta(SMeterMeta *pMeta, SShowObj *pShow, void *pConn) {
return
0
;
return
0
;
}
}
int32_t
mgmtRetrieveSuperTables
(
SShowObj
*
pShow
,
char
*
data
,
int32_t
rows
,
void
*
pConn
)
{
int32_t
mgmtRetrieveS
howS
uperTables
(
SShowObj
*
pShow
,
char
*
data
,
int32_t
rows
,
void
*
pConn
)
{
int32_t
numOfRows
=
0
;
int32_t
numOfRows
=
0
;
// char * pWrite;
// char * pWrite;
// int32_t cols = 0;
// int32_t cols = 0;
...
@@ -602,3 +602,45 @@ int32_t mgmtGetTagsLength(SSuperTableObj* pSuperTable, int32_t col) { // length
...
@@ -602,3 +602,45 @@ int32_t mgmtGetTagsLength(SSuperTableObj* pSuperTable, int32_t col) { // length
return
len
;
return
len
;
}
}
int32_t
mgmtSetSchemaFromSuperTable
(
SSchema
*
pSchema
,
SSuperTableObj
*
pTable
)
{
int32_t
numOfCols
=
pTable
->
numOfColumns
+
pTable
->
numOfTags
;
for
(
int32_t
i
=
0
;
i
<
numOfCols
;
++
i
)
{
strcpy
(
pSchema
->
name
,
pTable
->
schema
[
i
].
name
);
pSchema
->
type
=
pTable
->
schema
[
i
].
type
;
pSchema
->
bytes
=
htons
(
pTable
->
schema
[
i
].
bytes
);
pSchema
->
colId
=
htons
(
pTable
->
schema
[
i
].
colId
);
pSchema
++
;
}
return
(
pTable
->
numOfColumns
+
pTable
->
numOfTags
)
*
sizeof
(
SSchema
);
}
int32_t
mgmtGetSuperTableMeta
(
SDbObj
*
pDb
,
SSuperTableObj
*
pTable
,
SMeterMeta
*
pMeta
,
bool
usePublicIp
)
{
pMeta
->
uid
=
htobe64
(
pTable
->
uid
);
pMeta
->
sid
=
htonl
(
pTable
->
sid
);
pMeta
->
vgid
=
htonl
(
pTable
->
vgId
);
pMeta
->
sversion
=
htons
(
pTable
->
sversion
);
pMeta
->
precision
=
pDb
->
cfg
.
precision
;
pMeta
->
numOfTags
=
pTable
->
numOfTags
;
pMeta
->
numOfColumns
=
htons
(
pTable
->
numOfColumns
);
pMeta
->
tableType
=
pTable
->
type
;
pMeta
->
contLen
=
sizeof
(
SMeterMeta
)
+
mgmtSetSchemaFromSuperTable
(
pMeta
->
schema
,
pTable
);
SVgObj
*
pVgroup
=
mgmtGetVgroup
(
pTable
->
vgId
);
if
(
pVgroup
==
NULL
)
{
return
TSDB_CODE_INVALID_TABLE
;
}
for
(
int32_t
i
=
0
;
i
<
TSDB_VNODES_SUPPORT
;
++
i
)
{
if
(
usePublicIp
)
{
pMeta
->
vpeerDesc
[
i
].
ip
=
pVgroup
->
vnodeGid
[
i
].
publicIp
;
pMeta
->
vpeerDesc
[
i
].
vnode
=
htonl
(
pVgroup
->
vnodeGid
[
i
].
vnode
);
}
else
{
pMeta
->
vpeerDesc
[
i
].
ip
=
pVgroup
->
vnodeGid
[
i
].
ip
;
pMeta
->
vpeerDesc
[
i
].
vnode
=
htonl
(
pVgroup
->
vnodeGid
[
i
].
vnode
);
}
}
return
TSDB_CODE_SUCCESS
;
}
src/mnode/src/mgmtTable.c
浏览文件 @
03da5549
...
@@ -99,6 +99,24 @@ STableInfo* mgmtGetTableByPos(uint32_t dnodeIp, int32_t vnode, int32_t sid) {
...
@@ -99,6 +99,24 @@ STableInfo* mgmtGetTableByPos(uint32_t dnodeIp, int32_t vnode, int32_t sid) {
return
NULL
;
return
NULL
;
}
}
int32_t
mgmtGetTableMeta
(
SDbObj
*
pDb
,
STableInfo
*
pTable
,
SMeterMeta
*
pMeta
,
bool
usePublicIp
)
{
if
(
pTable
->
type
==
TSDB_TABLE_TYPE_CHILD_TABLE
)
{
mgmtGetChildTableMeta
(
pDb
,
(
SChildTableObj
*
)
pTable
,
pMeta
,
usePublicIp
);
}
else
if
(
pTable
->
type
==
TSDB_TABLE_TYPE_STREAM_TABLE
)
{
mgmtGetStreamTableMeta
(
pDb
,
(
SStreamTableObj
*
)
pTable
,
pMeta
,
usePublicIp
);
}
else
if
(
pTable
->
type
==
TSDB_TABLE_TYPE_NORMAL_TABLE
)
{
mgmtGetNormalTableMeta
(
pDb
,
(
SNormalTableObj
*
)
pTable
,
pMeta
,
usePublicIp
);
}
else
if
(
pTable
->
type
==
TSDB_TABLE_TYPE_SUPER_TABLE
)
{
mgmtGetSuperTableMeta
(
pDb
,
(
SSuperTableObj
*
)
pTable
,
pMeta
,
usePublicIp
);
}
else
{
mTrace
(
"%s, uid:%"
PRIu64
" table meta retrieve failed, invalid type"
,
pTable
->
tableId
,
pTable
->
uid
);
return
TSDB_CODE_INVALID_TABLE
;
}
mTrace
(
"%s, uid:%"
PRIu64
" table meta is retrieved"
,
pTable
->
tableId
,
pTable
->
uid
);
return
TSDB_CODE_SUCCESS
;
}
int32_t
mgmtCreateTable
(
SDbObj
*
pDb
,
SCreateTableMsg
*
pCreate
)
{
int32_t
mgmtCreateTable
(
SDbObj
*
pDb
,
SCreateTableMsg
*
pCreate
)
{
STableInfo
*
pTable
=
mgmtGetTable
(
pCreate
->
tableId
);
STableInfo
*
pTable
=
mgmtGetTable
(
pCreate
->
tableId
);
if
(
pTable
!=
NULL
)
{
if
(
pTable
!=
NULL
)
{
...
@@ -229,7 +247,7 @@ void mgmtCleanUpMeters() {
...
@@ -229,7 +247,7 @@ void mgmtCleanUpMeters() {
mgmtCleanUpSuperTables
();
mgmtCleanUpSuperTables
();
}
}
int32_t
mgmtGetTableMeta
(
SMeterMeta
*
pMeta
,
SShowObj
*
pShow
,
void
*
pConn
)
{
int32_t
mgmtGet
Show
TableMeta
(
SMeterMeta
*
pMeta
,
SShowObj
*
pShow
,
void
*
pConn
)
{
// int32_t cols = 0;
// int32_t cols = 0;
//
//
// SDbObj *pDb = NULL;
// SDbObj *pDb = NULL;
...
@@ -292,7 +310,7 @@ static void mgmtVacuumResult(char *data, int32_t numOfCols, int32_t rows, int32_
...
@@ -292,7 +310,7 @@ static void mgmtVacuumResult(char *data, int32_t numOfCols, int32_t rows, int32_
}
}
}
}
int32_t
mgmtRetrieveTables
(
SShowObj
*
pShow
,
char
*
data
,
int32_t
rows
,
void
*
pConn
)
{
int32_t
mgmtRetrieve
Show
Tables
(
SShowObj
*
pShow
,
char
*
data
,
int32_t
rows
,
void
*
pConn
)
{
int32_t
numOfRows
=
0
;
int32_t
numOfRows
=
0
;
// int32_t numOfRead = 0;
// int32_t numOfRead = 0;
// int32_t cols = 0;
// int32_t cols = 0;
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录