Ceph Crimson 网络调用和消息派发到OSD Shard分析

Ceph Crimson 网络调用和消息派发到OSD Shard分析

概述

本文档分析Ceph Crimson中网络层如何被调用,以及消息如何从网络层派发到OSD的shard的完整流程。

1. 网络层架构

1.1 核心组件

  • SocketMessenger: 消息传递器,管理连接和消息路由
  • SocketConnection: 表示一个网络连接
  • IOHandler: 处理连接的I/O操作(读/写消息)
  • ProtocolV2: 实现Ceph消息协议V2版本
  • FrameAssemblerV2: 负责消息帧的组装和解析
  • ChainedDispatchers: 链式分发器,将消息分发给多个Dispatcher

1.2 Shard管理

Crimson使用Seastar的shard模型,每个shard是一个独立的执行上下文(通常对应一个CPU核心)。网络连接可以在不同的shard之间迁移。

2. 网络监听和连接建立

2.1 启动监听

1
2
3
SocketMessenger::start(
const dispatchers_t& _dispatchers) {
assert(seastar::this_shard_id() == sid);

流程:

  1. SocketMessenger::start() 被调用,传入dispatchers列表
  2. 如果已绑定地址,调用 ShardedServerSocket::accept() 开始接受连接
  3. ShardedServerSocket 在所有shard上创建监听器(如果dispatch_only_on_this_shard=false

2.2 接受连接

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
seastar::future<>
ShardedServerSocket::accept(accept_func_t &&_fn_accept)
{
ceph_assert_always(seastar::this_shard_id() == primary_sid);
logger().debug("ShardedServerSocket({})::accept()...", listen_addr);
return this->container().invoke_on_all([_fn_accept](auto &ss) {
assert(ss.listener);
ss.fn_accept = _fn_accept;
// gate accepting
// ShardedServerSocket::shutdown() will drain the continuations in the gate
// so ignore the returned future
std::ignore = seastar::with_gate(ss.shutdown_gate, [&ss] {
return seastar::keep_doing([&ss] {
return ss.listener->accept(
).then([&ss](seastar::accept_result accept_result) {
#ifndef NDEBUG
if (ss.dispatch_only_on_primary_sid) {
// see seastar::listen_options::set_fixed_cpu()
ceph_assert_always(seastar::this_shard_id() == ss.primary_sid);
}
#endif
auto [socket, paddr] = std::move(accept_result);
entity_addr_t peer_addr;
peer_addr.set_sockaddr(&paddr.as_posix_sockaddr());
peer_addr.set_type(ss.listen_addr.get_type());
SocketRef _socket = std::make_unique<Socket>(
std::move(socket), Socket::side_t::acceptor,
peer_addr.get_port(), Socket::construct_tag{});
logger().debug("ShardedServerSocket({})::accept(): accepted peer {}, "
"socket {}, dispatch_only_on_primary_sid = {}",
ss.listen_addr, peer_addr, fmt::ptr(_socket.get()),
ss.dispatch_only_on_primary_sid);
std::ignore = seastar::with_gate(
ss.shutdown_gate,
[socket=std::move(_socket), peer_addr, &ss]() mutable {
return ss.fn_accept(std::move(socket), peer_addr
).handle_exception([&ss, peer_addr](auto eptr) {
const char *e_what;
try {
std::rethrow_exception(eptr);
} catch (std::exception &e) {
e_what = e.what();
}
logger().error("ShardedServerSocket({})::accept(): "
"fn_accept(s, {}) got unexpected exception {}",
ss.listen_addr, peer_addr, e_what);
ceph_abort();
});
});
});
}).handle_exception_type([&ss](const std::system_error& e) {
if (e.code() == std::errc::connection_aborted ||
e.code() == std::errc::invalid_argument) {
logger().debug("ShardedServerSocket({})::accept(): stopped ({})",
ss.listen_addr, e.what());
} else {
throw;
}
}).handle_exception([&ss](auto eptr) {
const char *e_what;
try {
std::rethrow_exception(eptr);
} catch (std::exception &e) {
e_what = e.what();
}
logger().error("ShardedServerSocket({})::accept(): "
"got unexpected exception {}", ss.listen_addr, e_what);
ceph_abort();
});
});
});
}

当新连接到达时:

  1. ShardedServerSocket::accept() 在循环中接受连接
  2. 创建 Socket 对象
  3. 调用 fn_accept 回调(由 SocketMessenger::accept() 提供)
  4. SocketMessenger::accept() 创建 SocketConnection 并开始握手

3. 消息读取流程

3.1 启动消息读取

当连接建立并完成握手后,IOHandler 开始读取消息:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
void IOHandler::do_in_dispatch()
{
shard_states->enter_in_dispatching();
shard_states->dispatch_in_background(
"do_in_dispatch", conn, [this, &ctx=*shard_states] {
return seastar::keep_doing([this, &ctx] {
return frame_assembler->read_main_preamble<false>(
).then([this, &ctx](auto ret) {
switch (ret.tag) {
case Tag::MESSAGE: {
size_t msg_size = get_msg_size(*ret.rx_frame_asm);
return seastar::futurize_invoke([this] {
// throttle_message() logic
if (!conn.policy.throttler_messages) {
return seastar::now();
}
// TODO: message throttler
ceph_abort_msg("TODO");
return seastar::now();
}).then([this, msg_size] {
// throttle_bytes() logic
if (!conn.policy.throttler_bytes) {
return seastar::now();
}
if (!msg_size) {
return seastar::now();
}
logger().trace("{} wants {} bytes from policy throttler {}/{}",
conn, msg_size,
conn.policy.throttler_bytes->get_current(),
conn.policy.throttler_bytes->get_max());
return conn.policy.throttler_bytes->get(msg_size);
}).then([this, msg_size, &ctx] {
// TODO: throttle_dispatch_queue() logic
utime_t throttle_stamp{seastar::lowres_system_clock::now()};
return read_message(ctx, throttle_stamp, msg_size);
});
}
case Tag::ACK:
return frame_assembler->read_frame_payload<false>(
).then([this](auto payload) {
// handle_message_ack() logic
auto ack = AckFrame::Decode(payload->back());
logger().debug("{} GOT AckFrame: seq={}", conn, ack.seq());
ack_out_sent(ack.seq());
});
case Tag::KEEPALIVE2:
return frame_assembler->read_frame_payload<false>(
).then([this](auto payload) {
// handle_keepalive2() logic
auto keepalive_frame = KeepAliveFrame::Decode(payload->back());
logger().debug("{} GOT KeepAliveFrame: timestamp={}",
conn, keepalive_frame.timestamp());
// notify keepalive ack
next_keepalive_ack = keepalive_frame.timestamp();
if (seastar::this_shard_id() == get_shard_id()) {
notify_out_dispatch();
}

last_keepalive = seastar::lowres_system_clock::now();
});
case Tag::KEEPALIVE2_ACK:
return frame_assembler->read_frame_payload<false>(
).then([this](auto payload) {
// handle_keepalive2_ack() logic
auto keepalive_ack_frame = KeepAliveFrameAck::Decode(payload->back());
auto _last_keepalive_ack =
seastar::lowres_system_clock::time_point{keepalive_ack_frame.timestamp()};
set_last_keepalive_ack(_last_keepalive_ack);
logger().debug("{} GOT KeepAliveFrameAck: timestamp={}",
conn, _last_keepalive_ack);
});
default: {
logger().warn("{} do_in_dispatch() received unexpected tag: {}",
conn, static_cast<uint32_t>(ret.tag));
abort_in_fault();
}
}
});
}).handle_exception([this, &ctx](std::exception_ptr eptr) {
const char *e_what;
try {
std::rethrow_exception(eptr);
} catch (std::exception &e) {
e_what = e.what();
}

auto io_state = ctx.get_io_state();
if (io_state == io_state_t::open) {
auto cc_seq = proto_crosscore.prepare_submit();
logger().info("{} do_in_dispatch(): fault at {}, {}, going to delay -- {}, "
"send {} notify_out_fault()",
conn, io_state, io_stat_printer{*this}, e_what, cc_seq);
do_set_io_state(io_state_t::delay);
shard_states->dispatch_in_background(
"notify_out_fault(in)", conn, [this, cc_seq, eptr] {
auto states = get_states();
return seastar::smp::submit_to(
conn.get_messenger_shard_id(), [this, cc_seq, eptr, states] {
return handshake_listener->notify_out_fault(
cc_seq, "do_in_dispatch", eptr, states);
});
});
} else {
if (io_state != io_state_t::switched) {
logger().info("{} do_in_dispatch(): fault at {}, {} -- {}",
conn, io_state, io_stat_printer{*this}, e_what);
} else {
logger().info("{} do_in_dispatch(): fault at {} -- {}",
conn, io_state, e_what);
}
}
}).finally([&ctx] {
ctx.exit_in_dispatching();
});
});
}

流程:

  1. do_in_dispatch() 进入消息读取循环
  2. 调用 frame_assembler->read_main_preamble() 读取消息头
  3. 根据消息类型(MESSAGE、ACK、KEEPALIVE等)进行不同处理
  4. 对于MESSAGE类型,调用 read_message() 读取完整消息

3.2 读取和解析消息

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
seastar::future<>
IOHandler::read_message(
shard_states_t &ctx,
utime_t throttle_stamp,
std::size_t msg_size)
{
return frame_assembler->read_frame_payload<false>(
).then([this, throttle_stamp, msg_size, &ctx](auto payload) {
if (unlikely(ctx.get_io_state() != io_state_t::open)) {
logger().debug("{} triggered {} during read_message()",
conn, ctx.get_io_state());
abort_protocol();
}

utime_t recv_stamp{seastar::lowres_system_clock::now()};

// we need to get the size before std::moving segments data
auto msg_frame = MessageFrame::Decode(*payload);
// XXX: paranoid copy just to avoid oops
ceph_msg_header2 current_header = msg_frame.header();

logger().trace("{} got {} + {} + {} byte message,"
" envelope type={} src={} off={} seq={}",
conn,
msg_frame.front_len(),
msg_frame.middle_len(),
msg_frame.data_len(),
(uint16_t)current_header.type,
conn.get_peer_name(),
(uint16_t)current_header.data_off,
(uint32_t)current_header.seq);

ceph_msg_header header{current_header.seq,
current_header.tid,
current_header.type,
current_header.priority,
current_header.version,
ceph_le32(msg_frame.front_len()),
ceph_le32(msg_frame.middle_len()),
ceph_le32(msg_frame.data_len()),
current_header.data_off,
conn.get_peer_name(),
current_header.compat_version,
current_header.reserved,
ceph_le32(0)};
ceph_msg_footer footer{ceph_le32(0), ceph_le32(0),
ceph_le32(0), ceph_le64(0), current_header.flags};

Message *message = decode_message(nullptr, 0, header, footer,
msg_frame.front(), msg_frame.middle(), msg_frame.data(), nullptr);
if (!message) {
logger().warn("{} decode message failed", conn);
abort_in_fault();
}

// store reservation size in message, so we don't get confused
// by messages entering the dispatch queue through other paths.
message->set_dispatch_throttle_size(msg_size);

message->set_throttle_stamp(throttle_stamp);
message->set_recv_stamp(recv_stamp);
message->set_recv_complete_stamp(utime_t{seastar::lowres_system_clock::now()});

// check received seq#. if it is old, drop the message.
// note that incoming messages may skip ahead. this is convenient for the
// client side queueing because messages can't be renumbered, but the (kernel)
// client will occasionally pull a message out of the sent queue to send
// elsewhere. in that case it doesn't matter if we "got" it or not.
uint64_t cur_seq = in_seq;
if (message->get_seq() <= cur_seq) {
logger().error("{} got old message {} <= {} {}, discarding",
conn, message->get_seq(), cur_seq, *message);
if (HAVE_FEATURE(conn.features, RECONNECT_SEQ) &&
local_conf()->ms_die_on_old_message) {
ceph_assert(0 == "old msgs despite reconnect_seq feature");
}
return seastar::now();
} else if (message->get_seq() > cur_seq + 1) {
logger().error("{} missed message? skipped from seq {} to {}",
conn, cur_seq, message->get_seq());
if (local_conf()->ms_die_on_skipped_message) {
ceph_assert(0 == "skipped incoming seq");
}
}

// note last received message.
in_seq = message->get_seq();
if (conn.policy.lossy) {
logger().debug("{} <== #{} === {} ({})",
conn,
message->get_seq(),
*message,
message->get_type());
} else {
logger().debug("{} <== #{},{} === {} ({})",
conn,
message->get_seq(),
(uint32_t)current_header.ack_seq,
*message,
*message,
message->get_type());
}

// notify ack
if (!conn.policy.lossy) {
++ack_left;
notify_out_dispatch();
}

ack_out_sent(current_header.ack_seq);

// TODO: change MessageRef with seastar::shared_ptr
auto msg_ref = MessageRef{message, false};
assert(ctx.get_io_state() == io_state_t::open);
assert(get_io_state() == io_state_t::open);
ceph_assert_always(conn_ref);

// throttle the reading process by the returned future
return dispatchers.ms_dispatch(conn_ref, std::move(msg_ref));
// user can make changes
});
}

关键步骤:

  1. 读取消息帧的payload
  2. 解码消息头(MessageFrame::Decode
  3. 调用 decode_message() 创建 Message 对象
  4. 验证消息序列号
  5. 调用 dispatchers.ms_dispatch() 派发消息

4. 消息派发到Dispatcher

4.1 ChainedDispatchers派发

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
seastar::future<>
ChainedDispatchers::ms_dispatch(ConnectionRef conn,
MessageRef m) {
try {
for (auto& dispatcher : dispatchers) {
auto dispatched = dispatcher->ms_dispatch(conn, m);
if (dispatched.has_value()) {
return std::move(*dispatched
).handle_exception([conn] (std::exception_ptr eptr) {
logger().error("{} got unexpected exception in ms_dispatch() throttling {}",
*conn, eptr);
ceph_abort();
});
}
}
} catch (...) {
logger().error("{} got unexpected exception in ms_dispatch() {}",
*conn, std::current_exception());
ceph_abort();
}
if (!dispatchers.empty()) {
logger().error("ms_dispatch unhandled message {}", *m);
}
return seastar::now();
}

ChainedDispatchers 按顺序遍历所有dispatchers:

  • 如果某个dispatcher返回非空的future,表示它处理了该消息,停止遍历
  • 如果所有dispatchers都不处理,记录错误

4.2 OSD的Dispatcher实现

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
std::optional<seastar::future<>>
OSD::ms_dispatch(crimson::net::ConnectionRef conn, MessageRef m)
{
if (pg_shard_manager.is_stopping()) {
return seastar::now();
}
auto maybe_ret = do_ms_dispatch(conn, std::move(m));
if (!maybe_ret.has_value()) {
return std::nullopt;
}

gate.dispatch_in_background(
__func__, *this, [ret=std::move(maybe_ret.value())]() mutable {
return std::move(ret);
});
return seastar::now();
}

OSD的 ms_dispatch() 实现:

  1. 检查OSD是否正在停止
  2. 调用 do_ms_dispatch() 处理消息
  3. 如果返回future,在后台执行

4.3 跨Shard消息处理

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
std::optional<seastar::future<>>
OSD::do_ms_dispatch(
crimson::net::ConnectionRef conn,
MessageRef m)
{
if (seastar::this_shard_id() != PRIMARY_CORE) {
switch (m->get_type()) {
case CEPH_MSG_OSD_MAP:
case MSG_COMMAND:
case MSG_OSD_MARK_ME_DOWN:
// FIXME: order is not guaranteed in this path
return conn.get_foreign(
).then([this, m=std::move(m)](auto f_conn) {
return seastar::smp::submit_to(PRIMARY_CORE,
[f_conn=std::move(f_conn), m=std::move(m), this]() mutable {
auto conn = make_local_shared_foreign(std::move(f_conn));
auto ret = do_ms_dispatch(conn, std::move(m));
assert(ret.has_value());
return std::move(ret.value());
});
});
}
}

switch (m->get_type()) {

重要机制:

  • 如果消息不在PRIMARY_CORE上,某些消息类型(如OSD_MAP)会被转发到PRIMARY_CORE
  • 使用 seastar::smp::submit_to() 跨shard提交任务

5. Shard迁移机制

5.1 连接迁移到新Shard

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
seastar::future<>
IOHandler::to_new_sid(
cc_seq_t cc_seq,
seastar::shard_id new_sid,
ConnectionFRef conn_fref,
std::optional<bool> is_replace)
{
ceph_assert_always(seastar::this_shard_id() == get_shard_id());
if (!proto_crosscore.proceed_or_wait(cc_seq)) {
logger().debug("{} got {} to_new_sid(), wait at {}",
conn, cc_seq, proto_crosscore.get_in_seq());
return proto_crosscore.wait(cc_seq
).then([this, cc_seq, new_sid, is_replace,
conn_fref=std::move(conn_fref)]() mutable {
return to_new_sid(cc_seq, new_sid, std::move(conn_fref), is_replace);
});
}

bool is_accept_or_connect = is_replace.has_value();
logger().debug("{} got {} to_new_sid_1(new_sid={}, {}) at {}",
conn, cc_seq, new_sid,
fmt::format("{}",
is_accept_or_connect ?
(*is_replace ? "accept(replace)" : "accept(!replace)") :
"connect"),
io_stat_printer{*this});
auto next_cc_seq = ++cc_seq;

if (get_io_state() != io_state_t::drop) {
ceph_assert_always(conn_ref);
if (new_sid != seastar::this_shard_id()) {
dispatchers.ms_handle_shard_change(conn_ref, new_sid, is_accept_or_connect);
// user can make changes
}
} else {
// it is possible that both io_handler and protocolv2 are
// trying to close each other from different cores simultaneously.
assert(!protocol_is_connected);
}

if (get_io_state() != io_state_t::drop) {
if (is_accept_or_connect) {
// protocol_is_connected can be from true to true here if the replacing is
// happening to a connected connection.
} else {
ceph_assert_always(protocol_is_connected == false);
}
protocol_is_connected = true;
} else {
assert(!protocol_is_connected);
}

bool is_dropped = false;
if (get_io_state() == io_state_t::drop) {
is_dropped = true;
}
ceph_assert_always(get_io_state() != io_state_t::open);

// apply the switching atomically
ceph_assert_always(conn_ref);
conn_ref.reset();
auto prv_sid = get_shard_id();
ceph_assert_always(maybe_prv_shard_states == nullptr);
maybe_prv_shard_states = std::move(shard_states);
shard_states = shard_states_t::create_from_previous(
*maybe_prv_shard_states, new_sid);
assert(new_sid == get_shard_id());
// broadcast shard change to all the io waiters, atomically.
io_crosscore.reset_wait();

return seastar::smp::submit_to(new_sid,
[this, next_cc_seq, is_dropped, prv_sid, is_replace, conn_fref=std::move(conn_fref)]() mutable {
logger().debug("{} got {} to_new_sid_2(prv_sid={}, is_dropped={}, {}) at {}",
conn, next_cc_seq, prv_sid, is_dropped,
fmt::format("{}",
is_replace.has_value() ?
(*is_replace ? "accept(replace)" : "accept(!replace)") :
"connect"),
io_stat_printer{*this});

ceph_assert_always(seastar::this_shard_id() == get_shard_id());
ceph_assert_always(get_io_state() != io_state_t::open);
ceph_assert_always(!maybe_dropped_sid.has_value());
ceph_assert_always(proto_crosscore.proceed_or_wait(next_cc_seq));

if (is_dropped) {
ceph_assert_always(get_io_state() == io_state_t::drop);
ceph_assert_always(shard_states->assert_closed_and_exit());
maybe_dropped_sid = prv_sid;
// cleanup_prv_shard() will be done in a follow-up close_io()
} else {
// possible at io_state_t::drop

// previous shard is not cleaned,
// but close_io() is responsible to clean up the current shard,
// so cleanup the previous shard here.
shard_states->dispatch_in_background(
"cleanup_prv_sid", conn, [this, prv_sid] {
return cleanup_prv_shard(prv_sid);
});
maybe_notify_out_dispatch();
}

ceph_assert_always(!conn_ref);
// assign even if already dropping
conn_ref = make_local_shared_foreign(std::move(conn_fref));

if (get_io_state() != io_state_t::drop) {
if (is_replace.has_value()) {
dispatchers.ms_handle_accept(conn_ref, prv_sid, *is_replace);
} else {
dispatchers.ms_handle_connect(conn_ref, prv_sid);
}
// user can make changes
}
});
}

Shard迁移流程:

  1. 在旧shard上:保存当前状态,通知dispatchers shard变化
  2. 创建新shard的状态对象
  3. 提交任务到新shard
  4. 在新shard上:恢复连接引用,通知dispatchers连接已迁移

6. 消息发送流程

6.1 发送消息

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
seastar::future<> IOHandler::send(MessageURef _msg)
{
// may be invoked from any core
MessageFRef msg = seastar::make_foreign(std::move(_msg));
auto cc_seq = io_crosscore.prepare_submit();
auto source_core = seastar::this_shard_id();
logger().debug("{} send {} send() to core {} -- {}",
conn, cc_seq, get_shard_id(), *msg);
// sid may be changed on-the-fly during the submission
if (source_core == get_shard_id()) {
return do_send(cc_seq, source_core, std::move(msg));
} else {
return seastar::smp::submit_to(
get_shard_id(),
[this, cc_seq, source_core, msg=std::move(msg)]() mutable {
return send_recheck_shard(cc_seq, source_core, std::move(msg));
});
}
}

发送流程:

  1. 如果调用者不在连接的shard上,使用 smp::submit_to() 提交到正确的shard
  2. 调用 do_send() 将消息加入发送队列
  3. 通知输出调度器开始发送

7. 关键数据结构

7.1 Shard States

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
class shard_states_t {
public:
shard_states_t(seastar::shard_id _sid, io_state_t state)
: sid{_sid}, io_state{state}, gate{_sid} {}

seastar::shard_id get_shard_id() const {
return sid;
}

io_state_t get_io_state() const {
assert(seastar::this_shard_id() == sid);
return io_state;
}

void set_io_state(io_state_t new_state) {
assert(seastar::this_shard_id() == sid);
assert(io_state != new_state);
pr_io_state_changed.set_value();
pr_io_state_changed = seastar::promise<>();
if (io_state == io_state_t::open) {
// from open
if (out_dispatching) {
ceph_assert_always(!out_exit_dispatching.has_value());
out_exit_dispatching = seastar::promise<>();
}
}
io_state = new_state;
}

seastar::future<> wait_state_change() {
assert(seastar::this_shard_id() == sid);
return pr_io_state_changed.get_future();
}

template <typename Func>
void dispatch_in_background(
const char *what, SocketConnection &who, Func &&func) {
assert(seastar::this_shard_id() == sid);
ceph_assert_always(!gate.is_closed());
gate.dispatch_in_background(what, who, std::move(func));
}

void enter_in_dispatching() {
assert(seastar::this_shard_id() == sid);
assert(io_state == io_state_t::open);
ceph_assert_always(!in_exit_dispatching.has_value());
in_exit_dispatching = seastar::promise<>();
}

void exit_in_dispatching() {
assert(seastar::this_shard_id() == sid);
assert(io_state != io_state_t::open);
ceph_assert_always(in_exit_dispatching.has_value());
in_exit_dispatching->set_value();
in_exit_dispatching = std::nullopt;
}

bool try_enter_out_dispatching(SocketConnection &conn) {
assert(seastar::this_shard_id() == sid);
if (out_dispatching) {
// already dispatching out
return false;
}
switch (io_state) {
case io_state_t::open:
[[fallthrough]];
case io_state_t::delay:
out_dispatching = true;
return true;
case io_state_t::drop:
[[fallthrough]];
case io_state_t::switched:
// do not dispatch out
return false;
default:
crimson::get_logger(ceph_subsys_ms).error(
"{} try_enter_out_dispatching() got wrong io_state {}",
conn, io_state);
ceph_abort_msg("impossible");
}
}

void notify_out_dispatching_stopped(
const char *what, SocketConnection &conn);

void exit_out_dispatching(
const char *what, SocketConnection &conn) {
assert(seastar::this_shard_id() == sid);
ceph_assert_always(out_dispatching);
out_dispatching = false;
notify_out_dispatching_stopped(what, conn);
}

seastar::future<> wait_io_exit_dispatching();

seastar::future<> close() {
assert(seastar::this_shard_id() == sid);
assert(!gate.is_closed());
return gate.close();
}

bool assert_closed_and_exit() const {
assert(seastar::this_shard_id() == sid);
if (gate.is_closed()) {
ceph_assert_always(io_state == io_state_t::drop ||
io_state == io_state_t::switched);
ceph_assert_always(!out_dispatching);
ceph_assert_always(!out_exit_dispatching);
ceph_assert_always(!in_exit_dispatching);
return true;
} else {
return false;
}
}

static shard_states_ref_t create(
seastar::shard_id sid, io_state_t state) {
return std::make_unique<shard_states_t>(sid, state);
}

static shard_states_ref_t create_from_previous(
shard_states_t &prv_states, seastar::shard_id new_sid);

private:
const seastar::shard_id sid;
io_state_t io_state;

crimson::common::Gated gate;
seastar::promise<> pr_io_state_changed;
bool out_dispatching = false;
std::optional<seastar::promise<>> out_exit_dispatching;
std::optional<seastar::promise<>> in_exit_dispatching;
};

shard_states_t 管理每个shard的状态:

  • io_state: delay/open/drop/switched
  • gate: 用于同步后台任务
  • out_dispatching/in_dispatching: 跟踪I/O调度状态

8. 总结

8.1 消息接收流程

  1. 网络层接受连接ShardedServerSocket::accept()
  2. 创建连接对象SocketMessenger::accept()SocketConnection
  3. 握手完成ProtocolV2 完成握手
  4. 启动消息读取IOHandler::do_in_dispatch()
  5. 读取消息帧FrameAssemblerV2::read_main_preamble()
  6. 解析消息IOHandler::read_message()
  7. 派发给DispatcherChainedDispatchers::ms_dispatch()
  8. OSD处理消息OSD::ms_dispatch()OSD::do_ms_dispatch()

8.2 Shard分配机制

  • 连接可以在不同shard之间迁移
  • 使用 seastar::smp::submit_to() 跨shard提交任务
  • 每个shard维护独立的状态(shard_states_t
  • OSD的某些消息必须在PRIMARY_CORE上处理

8.3 关键设计点

  1. 异步I/O: 所有网络操作都是异步的,使用Seastar的future/promise模型
  2. Shard隔离: 每个shard独立处理,避免锁竞争
  3. 状态管理: 使用状态机管理连接和I/O状态
  4. 错误处理: 完善的异常处理和连接恢复机制

文章互动

阅读 --

留言

0 条留言

正在加载留言…