Skip to content

feat: add common transport-level handshake for RDMA and UBSHM - #3535

Open
zchuango wants to merge 1 commit into
apache:masterfrom
LinQuickDev:transport-handshake
Open

zchuango wants to merge 1 commit into
apache:masterfrom
LinQuickDev:transport-handshake

Conversation

@zchuango

Copy link
Copy Markdown
Contributor

What problem does this PR solve?

Issue Number: N/A

Related Discussion: #3432

Problem Summary:

RDMA, UBSHM, and URMA use similar connection setup flows: establish a TCP control connection, exchange handshake messages, prepare transport-specific resources, negotiate whether the high-speed transport can be used, and fall back to TCP when necessary.

These handshake flows were previously implemented separately in individual transports, resulting in duplicated framing, TCP handshake I/O, state transitions, fallback handling, and connection orchestration.

As discussed in the related Discussion, the handshake orchestration should be moved above individual high-speed transports, while transport-specific resource management and data-plane operations remain inside each transport.

This PR introduces the common transport-level handshake framework and migrates both RDMA and UBSHM to it. URMA can be migrated to the same framework in a follow-up change.

What is changed and the side effects?

Changed:

  • Add AdapterTransport as the top-level transport for TCP, RDMA, and UBSHM sockets.
  • Add common handshake session, framing, I/O, and adapter abstractions.
  • Centralize handshake orchestration, state transitions, upgrade selection, and TCP fallback in the common transport layer.
  • Move RDMA handshake orchestration out of RdmaEndpoint.
  • Move UBSHM handshake orchestration out of UBShmEndpoint.
  • Keep transport-specific resource allocation, negotiation, activation, and data-plane operations inside their respective transports and endpoints.
  • Preserve the existing RDMA protocol ID and registered rdma_handshake protocol name.
  • Preserve the existing RDMA v2/v3 wire formats.
  • Preserve the existing UBSHM handshake wire format.
  • Preserve TCP fallback when the requested high-speed transport is unavailable or negotiation fails.
  • Add common handshake tests and transport-specific compatibility/fallback coverage for RDMA and UBSHM.

Side effects:

  • Performance effects:

    • No intended data-plane performance change.
    • The common handshake framework only affects connection setup and transport negotiation.
  • Breaking backward compatibility:

    • No intended breaking change.
    • Existing RDMA v2/v3 wire formats are preserved.
    • Existing UBSHM wire format is preserved.
    • Existing TCP fallback behavior is preserved.

Check List:

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

Unresolved critical and moderate handshake, transport-selection, fallback, initialization, and I/O issues remain.

Once you've addressed the issues Copilot identified, you can request another Copilot review.

Pull request overview

Adds a shared transport-level handshake framework for RDMA and UBSHM, centralizing framing, negotiation, upgrade selection, and TCP fallback.

Changes:

  • Adds common transport, session, framing, I/O, and adapter abstractions.
  • Migrates RDMA and UBSHM handshake orchestration.
  • Adds compatibility, fallback, handshake tests, and build integration.
File summaries
File Reviewed change / final finding
test/brpc_ubring_unittest.cpp UBSHM compatibility and fallback tests.
test/brpc_transport_handshake_unittest.cpp Common framing and handshake state-transition tests.
src/brpc/ubshm/ub_endpoint.h UBSHM endpoint interface migration.
src/brpc/ubshm_transport.h UBSHM transport declarations.
src/brpc/ubshm_transport.cpp UBSHM resource and fallback integration. Moderate (3 votes): fallback leaves resources active. Moderate (1 vote): allocation failure leaves a null endpoint on an upgrade-capable transport.
src/brpc/transport_handshake.h Shared handshake phases and session contracts. Critical (1 vote): relaxed phase publication can leave handshake reads blocked; release synchronization is needed.
src/brpc/transport_handshake.cpp Common handshake state transitions.
src/brpc/transport_factory.h Transport factory interface.
src/brpc/transport_factory.cpp Adapter transport construction.
src/brpc/socket.h Socket transport integration.
src/brpc/rdma/rdma_handshake.h RDMA handshake declarations.
src/brpc/rdma/rdma_handshake.cpp RDMA handshake implementation integration.
src/brpc/rdma/rdma_handshake_server.h RDMA server handshake declarations.
src/brpc/rdma/rdma_handshake_server.cpp RDMA server handshake implementation.
src/brpc/rdma/rdma_endpoint.h RDMA endpoint interface.
src/brpc/rdma_transport.h RDMA transport declarations.
src/brpc/rdma_transport.cpp RDMA resource and fallback integration. Moderate (3 votes): fallback does not release allocated resources. Moderate (1 vote): allocation failure leaves a null endpoint on an upgrade-capable transport.
src/brpc/rdma_handshake.proto RDMA handshake wire definitions.
src/brpc/policy/transport_handshake_protocol.h Common handshake protocol declarations.
src/brpc/policy/transport_handshake_protocol.cpp Common handshake parser dispatch.
src/brpc/policy/rdma_handshake_protocol.h RDMA compatibility API declarations.
src/brpc/policy/rdma_handshake_protocol.cpp RDMA compatibility facade.
src/brpc/input_messenger.h Handshake and input-messenger integration.
src/brpc/handshake/ubshm_handshake.h UBSHM handshake adapter declarations.
src/brpc/handshake/ubshm_handshake.cpp Critical (2 votes): upgrade selection permits wrong-mode transport casts. Moderate (1 vote): coalesced application bytes are rejected. Critical (1 vote): short UBSHM names can cause a 48-byte over-read.
src/brpc/handshake/rdma_handshake.h RDMA handshake adapter declarations.
src/brpc/handshake/rdma_handshake.cpp Critical (2 votes): upgrade selection permits wrong-mode transport casts. Moderate (1 vote): coalesced application bytes are rejected.
src/brpc/handshake/rdma_handshake_constants.h RDMA handshake wire-format constants.
src/brpc/handshake/handshake_io.h Handshake I/O interface.
src/brpc/handshake/handshake_io.cpp Moderate (2 votes): interrupted reads should retry on EINTR.
src/brpc/handshake/handshake_frame.h Handshake frame specifications.
src/brpc/handshake/handshake_frame.cpp Handshake frame encoding and parsing.
src/brpc/handshake/handshake_adapter.h Handshake adapter interfaces.
src/brpc/handshake/handshake_adapter.cpp Handshake/input-messenger bridge.
src/brpc/global.cpp Handshake protocol registration.
src/brpc/adapter_transport.h Adapter transport interface. Critical (1 vote): capability checks must be transport-mode-specific before concrete downcasts.
src/brpc/adapter_transport.cpp Common orchestration and fallback. Moderate (1 vote): post-upgrade UBSHM TCP data needs rejection handling. Moderate (1 vote): StopConnect must cancel blocked handshake work.
Makefile Build-source integration.
docs/cn/handshake_common_design.md Common handshake design documentation.
CMakeLists.txt Source and protobuf integration.
BUILD.bazel Bazel source integration.
Review details

Suppressed comments (8)

src/brpc/adapter_transport.cpp:340

  • For a server UBSHM socket this callback remains registered on the TCP control fd after the handshake, but it directly invokes InputMessenger::OnNewMessages in every phase. Once the upgrade is established, TCP is only the control channel and application data should arrive from UBRing; a later TCP write can otherwise be parsed and dispatched as an ordinary RPC (or re-enter the handshake parser) instead of being rejected. Use the same post-upgrade unexpected-TCP check as the RDMA server path while retaining normal parsing during TCP fallback.
            _on_edge_trigger = InputMessenger::OnNewMessages;

src/brpc/adapter_transport.cpp:62

  • StartConnect launches ProcessClientHandshake, whose task owns a SocketUniquePtr, but this StopConnect implementation is a no-op. If the socket is failed while the handshake is blocked in SocketHandshakeIO::ReadExact, nothing wakes the handshake's read butex or cancels the bthread; the task keeps the socket referenced, so recycling cannot reach StopConnect and the connection can leak a bthread/reference indefinitely. Keep a cancellable handshake handle (or make the handshake I/O observe failure and wake/abort) and cancel it here.
    void StopConnect(Socket*) override {}

src/brpc/handshake/handshake_io.cpp:117

  • write(2) is also allowed to return EINTR, but this loop immediately reports it as a fatal error. Retry interrupted writes before the EAGAIN wait path, otherwise a signal during the handshake can spuriously fail the connection.
        if (errno != EAGAIN) {
            return -1;
        }

src/brpc/handshake/rdma_handshake.cpp:551

  • As with UBSHM, the ACK parser may leave application bytes in source when the peer coalesces them with the handshake. This check turns that case into a failed RDMA connection even though StandardHandshakeAdapter is designed to return PARSE_ERROR_TRY_OTHERS and let the normal protocol parser consume the remaining buffer. Validate the transport state without rejecting buffered application data.
        if (!source->empty()) {
            return STEP_ERROR;
        }

src/brpc/handshake/ubshm_handshake.cpp:369

  • HandshakeSession::RunServer invokes validate_established() before set_high_speed_active(). Consequently UpgradeActive() is still false here on every otherwise valid UBSHM handshake, so this callback returns STEP_ERROR and the session never reaches ESTABLISHED. Remove this pre-activation check or move validation to a point where activation is already published.
        if (!source->empty() ||
            !transport->UpgradeActive()) {

src/brpc/handshake/ubshm_handshake.cpp:370

  • The ACK has already been consumed when this validation runs, but any application bytes coalesced in the same TCP read remain in source. Returning STEP_ERROR rejects that valid stream instead of returning TRY_OTHERS so InputMessengerProcessor can parse the remaining bytes. The validator should not require the input buffer to be empty.
        if (!source->empty() ||
            !transport->UpgradeActive()) {
            return STEP_ERROR;

src/brpc/rdma_transport.cpp:48

  • If new (std::nothrow) fails, this branch marks the socket failed but continues initialization with _rdma_ep == nullptr. RdmaTransport remains upgrade-capable, so later handshake/resource calls dereference the null endpoint (and GetRdmaEp() can hit CHECK). Propagate the initialization failure through the socket creation path or make the transport unavailable before returning.
    src/brpc/ubshm_transport.cpp:53
  • This allocation failure path records SetFailed but still returns from Init with _ub_ep == nullptr and a non-null high-speed transport. A subsequent handshake can call GetUBShmEp()/resource methods and dereference the missing endpoint. Propagate initialization failure or mark the transport unavailable before returning instead of leaving a partially initialized upgrade path.
  • Files reviewed: 43/44 changed files
  • Comments generated: 8
  • Review effort level: Lite

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread src/brpc/adapter_transport.h Outdated
Comment thread src/brpc/handshake/rdma_handshake.cpp Outdated
Comment thread src/brpc/handshake/ubshm_handshake.cpp Outdated
Comment thread src/brpc/handshake/ubshm_handshake.cpp Outdated
Comment thread src/brpc/transport_handshake.h Outdated
Comment thread src/brpc/handshake/handshake_io.cpp
Comment thread src/brpc/rdma_transport.cpp Outdated
Comment thread src/brpc/ubshm_transport.cpp Outdated

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Warning

Copilot couldn't run its full agentic review because it didn't start before the timeout. Make sure your repository has a runner available, or add a copilot-code-review.yml file specifying one with the runs-on attribute. See the docs for more details.

Pull request overview

Copilot reviewed 43 out of 44 changed files in this pull request and generated 1 comment.

Suppressed comments (3)

src/brpc/rdma_transport.cpp:1

  • Failing the Socket in RdmaTransport::Init when RdmaEndpoint allocation fails prevents the intended TCP fallback path from working (the connection becomes unusable even though the PR’s design explicitly supports fallback). Instead of socket->SetFailed(...) here, treat this as 'upgrade unavailable' (e.g., mark RDMA as OFF/disabled and let AdapterTransport continue with TCP), and ensure later RDMA-only operations are gated on upgrade readiness.
    src/brpc/ubshm_transport.cpp:1
  • Same issue as RDMA: marking the Socket failed on UBShmEndpoint allocation failure breaks TCP fallback (connection becomes unusable). This should be handled as 'UBSHM upgrade not available' and allow AdapterTransport to proceed with TCP, rather than failing the socket during init.
    src/brpc/ubshm/ub_endpoint.cpp:1
  • These UBSHM/UBRing flags have RDMA-specific descriptions ("RDMA polling mode" / "Disable bthread in RDMA"), which is misleading for operators. Update the strings to reference UBRing/UBSHM polling mode to match the flag names and file context.

Comment thread src/brpc/adapter_transport.cpp Outdated

@wwbmmm wwbmmm left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[replied by brpc-oncall robot] The handshake refactor looks reasonable and the latest revision addresses the earlier upgrade-selection, synchronization and resource-release concerns, but an accidental logic drop in RdmaEndpoint and the unrelated 2-space reformatting should be fixed.

Comment thread src/brpc/rdma/rdma_endpoint.cpp Outdated
int fd = _resource->comp_channel->fd;
GetGlobalEventDispatcher(fd, _socket->_io_event.bthread_tag()).RemoveConsumer(fd);
GetGlobalEventDispatcher(fd, _socket->_io_event.bthread_tag())
.RemoveConsumer(fd);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[replied by brpc-oncall robot] The inner if (_cq_sid != INVALID_SOCKET_ID) guard was dropped during the reformatting. RemoveConsumer(comp_channel->fd) is now called even when this endpoint never registered that fd with the event dispatcher, i.e. when DeallocateResources() runs after AllocateResources() but before StartCqEvents() (which is exactly what happens on TCP fallback via DeactivateUpgrade() -> RdmaEndpoint::Reset()). Calling epoll_ctl(EPOLL_CTL_DEL) on an fd that was never added fails with ENOENT and logs misleading "Fail to remove fd" warnings. Please restore the _cq_sid guard (and keep remove_consumer as before).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed. The _cq_sid != INVALID_SOCKET_ID guard has been restored before removing the RDMA completion-channel consumer.

Comment thread src/brpc/rdma_transport.cpp

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

Four unresolved review findings remain, including two critical issues.

Get a fresh assessment by requesting another Copilot review.

Review details

Suppressed comments (4)

src/brpc/adapter_transport.cpp:432

  • When endpoint allocation fails, upgrade_capable(_mode) is false, so this returns the raw TCP connector without running the common session and leaves the phase UNINITIALIZED. Client high-speed sockets use OnNewDataFromTcp, whose ProcessTcpEvent does not invoke InputMessenger in that phase, so TCP fallback responses are never parsed. Install the normal input callback for this no-upgrade client path (after the high-speed transport is initialized), or publish FALLBACK_TCP before returning.
std::shared_ptr<AppConnect> AdapterTransport::Connect() {
    if (upgrade_capable(_mode)) {
        return std::make_shared<AdapterConnect>(_default_connect);
    }
    return _tcp_transport->Connect();

src/brpc/adapter_transport.cpp:430

  • _app_connect is copied into SocketOptions.app_connect when the channel creates pooled/short sockets (socket.cpp:2805 and 2908). After the first high-speed socket, that value is already an AdapterConnect, so this branch wraps it in a second AdapterConnect; the outer handshake then invokes the inner one, causing a second handshake/resource allocation on an already-upgraded endpoint (RdmaEndpoint::DoAllocateResources asserts _resource == nullptr, and UBSHM has the analogous check). Return an existing adapter (and leave null server-side app-connects null) instead of nesting the wrapper.
std::shared_ptr<AppConnect> AdapterTransport::Connect() {
    if (upgrade_capable(_mode)) {
        return std::make_shared<AdapterConnect>(_default_connect);

src/brpc/handshake/rdma_handshake.cpp:560

  • The common driver invokes on_failed for send, ACK, and protocol errors after RDMA resources may already have been allocated, but this callback is a no-op. The failed session can therefore retain its QP/CQ resources until socket recycling; deactivate the upgrade here so the failed state also releases transport resources and stops transport-side activity.
    callbacks.transport.on_failed = []() {};

src/brpc/handshake/ubshm_handshake.cpp:385

  • The common driver invokes on_failed for send, ACK, and protocol errors after the UBRING resources and poller may already have been created, but this callback is a no-op. The failed session can retain the ring/poller until socket recycling; deactivate the upgrade here so failure immediately tears down transport-side resources.
    callbacks.transport.on_failed = []() {};
  • Files reviewed: 43/44 changed files
  • Comments generated: 2
  • Review effort level: Lite

Comment thread src/brpc/adapter_transport.cpp
Comment thread src/brpc/handshake/ubshm_handshake.cpp Outdated
Comment on lines +340 to +341
strncpy(remote_trx_shm.name, remote.shm_name,
SHM_MAX_NAME_BUFF_LEN);

@wwbmmm wwbmmm left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[replied by brpc-oncall robot] Copilot's two critical findings are valid and must be fixed (AdapterConnect must not be re-wrapped for pooled/short sockets, and the server must reject a non-NUL-terminated remote.shm_name before strncpy/strlen), but the client edge-trigger race claim does not apply to the sockets this PR creates.

Comment thread src/brpc/adapter_transport.cpp
@wwbmmm

wwbmmm commented Sep 22, 2026

Copy link
Copy Markdown
Contributor

please resolve conflicts

@wwbmmm wwbmmm left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This PR was intended to reduce duplicated code logic, but actually it added more lines of code than removed code.

Comment thread docs/cn/handshake_common_design.md Outdated
Comment thread test/brpc_rdma_unittest.cpp Outdated
@zchuango

Copy link
Copy Markdown
Contributor Author

@wwbmmm @chenBright The latest changes are ready and CI is green now. Could you please take another look when you have time? Thanks!

@wwbmmm

wwbmmm commented Sep 30, 2026

Copy link
Copy Markdown
Contributor

LGTM

const StepResult result = RunServerStep(source, socket);
if (result == STEP_NEED_MORE) {
if (GetSession(socket)->phase() == ACK_WAIT &&
socket->parsing_context() == NULL) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Use nullptr instead of NULL .

Comment thread src/brpc/global.cpp
Comment on lines 445 to +449
CONNECTION_TYPE_ALL, "rdma_handshake" };
if (RegisterProtocol(PROTOCOL_RDMA_HANDSHAKE, rdma_handshake_protocol) != 0) {
// Retain the existing enum value and registered name to avoid changing
// public protocol identifiers while widening the implementation from RDMA
// to all transport upgrades.
if (RegisterProtocol(PROTOCOL_RDMA_HANDSHAKE,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think modifying rdma_handshake and PROTOCOL_RDMA_HANDSHAKE is fine.

Comment thread src/brpc/transport_handshake.h Outdated
_phase.store(UNINITIALIZED, butil::memory_order_relaxed);
}

int phase(butil::memory_order order = butil::memory_order_acquire) const {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Require the caller to explicitly pass the memory_order instead of using a default argument.

Comment thread src/brpc/rdma_transport.cpp Outdated
extern SocketVarsCollector *g_vars;

RdmaTransport *RdmaTransport::Get(const Socket *socket) {
const AdapterTransport *adapter = AdapterTransport::Get(socket);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Same as above.

Comment thread src/brpc/rdma_transport.cpp Outdated

extern SocketVarsCollector *g_vars;

RdmaTransport *RdmaTransport::Get(const Socket *socket) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

RdmaTransport * -> RdmaTransport*

Comment thread src/brpc/adapter_transport.h Outdated
Comment on lines +38 to +40
friend class TransportFactory;
friend class RdmaTransport;
friend class UBShmTransport;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

No indentation required.


ParseResult ParseTransportHandshake(butil::IOBuf* source, Socket* socket,
bool /*read_eof*/, const void* /*arg*/) {
return AdapterTransport::Get(socket)->ProcessUpgradeReadable(source);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ParseTransportHandshake() unconditionally calls AdapterTransport::Get(socket), which performs a static_cast from the socket's transport.

TransportFactory still installs UrmaTransport for URMA sockets. When URMA falls back to TCP, its input reaches InputMessenger and the globally registered handshake parser can participate in protocol detection. This casts an unrelated object to AdapterTransport.

Even when the input does not match a handshake magic, ProcessUpgradeReadable() accesses handshake and connection-completion fields. These accesses produce undefined behavior and may crash.

Comment thread test/brpc_rdma_unittest.cpp Outdated
Comment on lines +69 to +75
extern ibv_cq* (*IbvCreateCq)(ibv_context*, int, void*, ibv_comp_channel*,
int);
extern int (*IbvDestroyCq)(ibv_cq*);
extern ibv_qp* (*IbvCreateQp)(ibv_pd*, ibv_qp_init_attr*);
extern int (*IbvModifyQp)(ibv_qp*, ibv_qp_attr*, ibv_qp_attr_mask);
extern int (*IbvQueryQp)(ibv_qp*, ibv_qp_attr*, ibv_qp_attr_mask, ibv_qp_init_attr*);
extern int (*IbvQueryQp)(ibv_qp*, ibv_qp_attr*, ibv_qp_attr_mask,
ibv_qp_init_attr*);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

No need to modify.

Comment on lines +383 to 406
sockaddr_in addr;
bzero((char*)&addr, sizeof(addr));
addr.sin_family = AF_INET;
addr.sin_port = htons(PORT);

butil::fd_guard sockfd(socket(AF_INET, SOCK_STREAM, 0));
ASSERT_TRUE(sockfd >= 0);
ASSERT_EQ(0, connect(sockfd, (sockaddr*)&addr, sizeof(sockaddr)));
usleep(100000); // wait for server to handle the msg
Socket* s = GetSocketFromServer(0);
ASSERT_EQ(handshake::UNINITIALIZED,
AdapterTransport::Get(s)->handshake_phase());

uint8_t data[rdma::HELLO_V2_MSG_LEN_MIN];
memcpy(data, "PRPC", 4); // send as normal baidu_std protocol
ASSERT_TRUE(WriteAll(sockfd, data, 4));
// Wait for the bytes to show up in the fd stream (baidu_std wants 12B of
// header, so they stay buffered). Waiting on the state instead would prove
// nothing: it is already UNINIT before the server has read anything.
ASSERT_TRUE(WaitForFdReadBuf(s, 4));
// A non-RDMA magic makes ParseRdmaHandshake return TRY_OTHERS and hand the
// bytes to other protocols; it does not touch the endpoint state, so it
// stays UNINIT (the old blocking handshake used to set FALLBACK_TCP here).
ASSERT_EQ(rdma::RdmaEndpoint::UNINIT, RdmaTransportOf(s)->_rdma_ep->_state);
ASSERT_EQ(4, write(sockfd, data, 4));
usleep(100000); // wait for server to handle the msg
// A non-RDMA magic makes the transport-handshake parser return TRY_OTHERS
// and hand the bytes to other protocols; it does not touch the endpoint
// state, so it stays UNINIT (the old blocking handshake used to set
// FALLBACK_TCP here).
ASSERT_EQ(handshake::UNINITIALIZED,
AdapterTransport::Get(s)->handshake_phase());

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

No need to modify.

Comment thread test/brpc_rdma_unittest.cpp Outdated
butil::fd_guard sockfd1(socket(AF_INET, SOCK_STREAM, 0));
ASSERT_TRUE(sockfd1 >= 0);
ASSERT_EQ(0, connect(sockfd1, (sockaddr*)&addr, sizeof(sockaddr)));
usleep(100000); // wait for server to handle the msg

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do not use sleep to ensure correctness. I think it's fine to stick with the original usage.

@zchuango

zchuango commented Oct 8, 2026 •

Copy link
Copy Markdown
Contributor Author

Thanks @chenBright for the detailed review!

I've updated this PR to address the recent feedback, including:

  • Fixing the URMA transport type handling.
  • Fixing UBSHM handshake continuation and shared-memory name handling.
  • Preventing TCP application data from being parsed after a successful upgrade.
  • Making handshake memory ordering explicit.
  • Improving resource cleanup and TCP fallback handling.
  • Adding regression tests for handshake fragmentation, TCP fallback, and resource cleanup.

I've also addressed the related style comments.
Could you please take another look when you have time?

Comment on lines +592 to +594
InputMessenger::OnNewMessagesUntil(socket, [](Socket* s) {
return Get(s)->handshake_phase() == handshake::ESTABLISHED;
});

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

UBSHM receive polling is registered during resource preparation, before the final TCP handshake ACK is consumed. UBShmEndpoint::PollIn() does not wait for handshake completion and feeds shared-memory bytes into the same fd_input_processor() used by the TCP handshake.

A client can write the final TCP ACK and immediately send its first RPC through shared memory. If the shared-memory poller runs before the server processes the TCP ACK, the retained handshake context interprets those RPC bytes as an ACK.

Comment on lines +381 to +388
sockaddr_in addr;
bzero((char*)&addr, sizeof(addr));
addr.sin_family = AF_INET;
addr.sin_port = htons(PORT);

butil::fd_guard sockfd(socket(AF_INET, SOCK_STREAM, 0));
ASSERT_TRUE(sockfd >= 0);
ASSERT_EQ(0, connect(sockfd, (sockaddr*)&addr, sizeof(sockaddr)));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are these changes necessary?

Comment thread test/brpc_rdma_unittest.cpp Outdated
Comment on lines +299 to +303
// observes a read event, not publication of the parser's mutable buffer.
static bool WaitForSocketRead(Socket* s, int64_t previous_read_us) {
return WaitUntil([s, previous_read_us] {
return s->_last_readtime_us.load(butil::memory_order_relaxed) !=
previous_read_us;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

WaitForSocketRead() observes _last_readtime_us, but ProcessNewMessage() updates that timestamp before invoking the parser.

When parsing the incomplete magic "RD", RunServer() temporarily publishes HELLO_WAIT and later restores UNINITIALIZED after determining that more bytes are needed.

@wwbmmm

wwbmmm commented Oct 8, 2026

Copy link
Copy Markdown
Contributor

please resolve merge conflict

@zchuango

zchuango commented Oct 8, 2026

Copy link
Copy Markdown
Contributor Author

please resolve merge conflict

ok👌

@chenBright

Copy link
Copy Markdown
Contributor

The method used to resolve conflicts seems problematic, resulting in numerous diffs involving the master branch.

I usually resolve conflicts by rebasing onto master and haven't encountered this issue.

@zchuango

zchuango commented Oct 9, 2026

Copy link
Copy Markdown
Contributor Author

The method used to resolve conflicts seems problematic, resulting in numerous diffs involving the master branch.

I usually resolve conflicts by rebasing onto master and haven't encountered this issue.

thanks for pointing this out, it was already fixed.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants