Skip to content
体验新版
项目
组织
正在加载...
登录
切换导航
打开侧边栏
oceanbase
oceanbase
提交
a97444e7
O
oceanbase
项目概览
oceanbase
/
oceanbase
大约 1 年 前同步成功
通知
261
Star
6084
Fork
1301
代码
文件
提交
分支
Tags
贡献者
分支图
Diff
Issue
0
列表
看板
标记
里程碑
合并请求
0
DevOps
流水线
流水线任务
计划
Wiki
0
Wiki
分析
仓库
DevOps
项目成员
Pages
O
oceanbase
项目概览
项目概览
详情
发布
仓库
仓库
文件
提交
分支
标签
贡献者
分支图
比较
Issue
0
Issue
0
列表
看板
标记
里程碑
合并请求
0
合并请求
0
Pages
DevOps
DevOps
流水线
流水线任务
计划
分析
分析
仓库分析
DevOps
Wiki
0
Wiki
成员
成员
收起侧边栏
关闭侧边栏
动态
分支图
创建新Issue
流水线任务
提交
Issue看板
体验新版 GitCode,发现更多精彩内容 >>
提交
a97444e7
编写于
10月 28, 2022
作者:
O
obdev
提交者:
wangzelin.wzl
10月 28, 2022
浏览文件
操作
浏览文件
下载
电子邮件补丁
差异文件
[server net keepalive] sync serialization size
上级
ab27c671
变更
1
显示空白变更内容
内联
并排
Showing
1 changed file
with
76 addition
and
4 deletion
+76
-4
deps/oblib/src/rpc/obrpc/ob_net_keepalive.cpp
deps/oblib/src/rpc/obrpc/ob_net_keepalive.cpp
+76
-4
未找到文件。
deps/oblib/src/rpc/obrpc/ob_net_keepalive.cpp
浏览文件 @
a97444e7
...
...
@@ -28,10 +28,13 @@
#include "lib/utility/ob_defer.h"
#include "lib/thread/ob_thread_name.h"
#include "lib/time/ob_time_utility.h"
#include "lib/utility/serialization.h"
#include "lib/utility/utility.h"
#include "rpc/frame/ob_net_easy.h"
#include "io/easy_negotiation.h"
using
namespace
oceanbase
::
common
;
using
namespace
oceanbase
::
common
::
serialization
;
using
namespace
oceanbase
::
lib
;
using
namespace
oceanbase
::
rpc
::
frame
;
namespace
oceanbase
...
...
@@ -44,6 +47,43 @@ namespace obrpc
#define WINDOW_MAX_FAILS 4 // 4 times
#define MAX_CREDIBLE_WINDOW 10 * 1000 * 1000 // 10s
constexpr
int32_t
KP_MAGIC
=
0x2c15c364
;
struct
Header
{
public:
Header
(
int32_t
data_len
=
0
)
:
magic_
(
KP_MAGIC
),
data_len_
(
data_len
)
{}
int
encode
(
char
*
buf
,
const
int64_t
buf_len
,
int64_t
&
pos
)
{
int
ret
=
OB_SUCCESS
;
if
(
OB_FAIL
(
encode_i32
(
buf
,
buf_len
,
pos
,
magic_
)))
{
_LOG_WARN
(
"encode magic failed, ret: %d, pos: %ld"
,
ret
,
pos
);
}
else
if
(
OB_FAIL
(
encode_i32
(
buf
,
buf_len
,
pos
,
data_len_
)))
{
_LOG_WARN
(
"encode data len failed, ret: %d, pos: %ld"
,
ret
,
pos
);
}
return
ret
;
}
int
decode
(
const
char
*
buf
,
const
int64_t
buf_len
,
int64_t
&
pos
)
{
int
ret
=
OB_SUCCESS
;
if
(
OB_FAIL
(
decode_i32
(
buf
,
buf_len
,
pos
,
&
magic_
)))
{
_LOG_WARN
(
"decode magic failed, ret: %d, pos: %ld"
,
ret
,
pos
);
}
else
if
(
magic_
!=
KP_MAGIC
)
{
ret
=
OB_ERR_UNEXPECTED
;
_LOG_WARN
(
"unexpected magic, magic: %d"
,
magic_
);
}
else
if
(
OB_FAIL
(
decode_i32
(
buf
,
buf_len
,
pos
,
&
data_len_
)))
{
_LOG_WARN
(
"decode data len failed, ret: %d, pos: %ld"
,
ret
,
pos
);
}
return
ret
;
}
int32_t
get_encoded_size
()
const
{
return
encoded_length_i32
(
magic_
)
+
encoded_length_i32
(
data_len_
);
}
int32_t
magic_
;
int32_t
data_len_
;
};
enum
{
UNCONNECT
=
0
,
CONNECTING
,
...
...
@@ -333,6 +373,7 @@ void ObNetKeepAlive::do_server_loop()
for
(
int
i
=
0
;
i
<
cnt
;
i
++
)
{
struct
server
*
s
=
(
struct
server
*
)
events
[
i
].
data
.
ptr
;
int
ev_fd
=
NULL
==
s
?
pipefd_
:
s
->
fd_
;
bool
need_disconn
=
false
;
if
(
NULL
==
s
)
{
struct
server
*
s
=
(
struct
server
*
)
ob_malloc
(
sizeof
(
struct
server
),
"KeepAliveServer"
);
if
(
NULL
==
s
)
{
...
...
@@ -383,10 +424,26 @@ void ObNetKeepAlive::do_server_loop()
char
data
=
PROTOCOL_DATA
;
while
((
n
=
read
(
ev_fd
,
&
data
,
sizeof
data
))
<
0
&&
errno
==
EINTR
);
if
(
n
<=
0
)
break
;
while
((
n
=
write
(
ev_fd
,
&
PROTOCOL_DATA
,
sizeof
PROTOCOL_DATA
))
<
0
&&
errno
==
EINTR
);
char
buf
[
128
];
const
int64_t
buf_len
=
sizeof
buf
;
int32_t
data_len
=
1
;
Header
header
(
data_len
);
int
tmp_ret
=
OB_SUCCESS
;
int64_t
pos
=
0
;
if
(
OB_SUCCESS
!=
(
tmp_ret
=
header
.
encode
(
buf
,
buf_len
,
pos
)))
{
_LOG_WARN
(
"encode failed, ret: %d, pos: %ld"
,
tmp_ret
,
pos
);
}
else
if
(
FALSE_IT
(
pos
+=
data_len
)
/*TODO: encode data*/
)
{
_LOG_WARN
(
"encode data failed: %d"
,
data_len
);
}
else
{
while
((
n
=
write
(
ev_fd
,
buf
,
pos
))
<
0
&&
errno
==
EINTR
);
need_disconn
=
n
<
pos
;
}
}
}
if
(
!
need_disconn
)
{
need_disconn
=
events
[
i
].
events
&
(
EPOLLRDHUP
|
EPOLLHUP
);
}
if
(
events
[
i
].
events
&
(
EPOLLRDHUP
|
EPOLLHUP
)
)
{
if
(
need_disconn
)
{
_LOG_INFO
(
"server connection closed, fd: %d, addr: %s"
,
ev_fd
,
NULL
==
s
?
""
:
addr_to_string
(
s
->
cli_addr_
));
epoll_ctl
(
epfd
,
EPOLL_CTL_DEL
,
ev_fd
,
NULL
);
close
(
ev_fd
);
...
...
@@ -464,9 +521,24 @@ void ObNetKeepAlive::do_client_loop()
_LOG_DEBUG
(
"update read ts, addr: %s, fd: %d, ts: %ld"
,
addr_to_string
(
rs
->
svr_addr_
),
ev_fd
,
rs
->
last_read_ts_
);
for
(;;)
{
ssize_t
n
=
-
1
;
char
data
=
PROTOCOL_DATA
;
while
((
n
=
read
(
ev_fd
,
&
data
,
sizeof
data
))
<
0
&&
errno
==
EINTR
);
char
buf
[
128
];
Header
header
;
int32_t
read_len
=
header
.
get_encoded_size
();
while
((
n
=
read
(
ev_fd
,
buf
,
read_len
))
<
0
&&
errno
==
EINTR
);
if
(
n
<=
0
)
break
;
int
tmp_ret
=
OB_SUCCESS
;
int64_t
pos
=
0
;
if
(
OB_SUCCESS
!=
(
tmp_ret
=
header
.
decode
(
buf
,
read_len
,
pos
)))
{
_LOG_WARN
(
"decode failed, ret: %d, pos: %ld"
,
tmp_ret
,
pos
);
}
else
{
char
data
[
512
];
// TODO
if
(
header
.
data_len_
>
sizeof
data
)
{
tmp_ret
=
OB_BUF_NOT_ENOUGH
;
_LOG_WARN
(
"data buf not enough: %d"
,
header
.
data_len_
);
}
else
{
while
((
n
=
read
(
ev_fd
,
data
,
header
.
data_len_
))
<
0
&&
errno
==
EINTR
);
}
}
do_rpin
(
client2rs
(
c
));
}
}
...
...
编辑
预览
Markdown
is supported
0%
请重试
或
添加新附件
.
添加附件
取消
You are about to add
0
people
to the discussion. Proceed with caution.
先完成此消息的编辑!
取消
想要评论请
注册
或
登录