Skip to content

GH-50282: [C++][FlightRPC] Refactor GRPC server and transport classes - #50407

Open
Alex-PLACET wants to merge 21 commits into
apache:mainfrom
Alex-PLACET:refactor/internal-transport
Open

Alex-PLACET wants to merge 21 commits into
apache:mainfrom
Alex-PLACET:refactor/internal-transport

Conversation

@Alex-PLACET

@Alex-PLACET Alex-PLACET commented Jul 7, 2026 •

Copy link
Copy Markdown
Contributor

The goal of the pullrequest is to refactor the FligtRPC code to prepare the async implementation. I moved code from grpc_server and transport_server to the _internal files which will be used in both sync and async server implmentations.
I also created dedicated function for specific responsabilities and reduce code complexity or duplication.

@github-actions github-actions Bot added the awaiting review Awaiting review label Jul 7, 2026
@github-actions

github-actions Bot commented Jul 7, 2026

Copy link
Copy Markdown

Thanks for opening a pull request!

If this is not a minor PR. Could you open an issue for this pull request on GitHub? https://github.com/apache/arrow/issues/new/choose

Opening GitHub issues ahead of time contributes to the Openness of the Apache Arrow project.

Then could you also rename the pull request title in the following format?

GH-${GITHUB_ISSUE_ID}: [${COMPONENT}] ${SUMMARY}

or

MINOR: [${COMPONENT}] ${SUMMARY}

See also:

@Alex-PLACET Alex-PLACET changed the title GH-#50282: [C++][FlightRPC] Refactor GRPC server and transport classes GH-50282: [C++][FlightRPC] Refactor GRPC server and transport classes Jul 7, 2026
@github-actions

github-actions Bot commented Jul 7, 2026

Copy link
Copy Markdown

⚠️ GitHub issue #50282 has been automatically assigned in GitHub to PR creator.

@Alex-PLACET

Copy link
Copy Markdown
Contributor Author

@raulcd Can you launch the CI please

@Alex-PLACET

Copy link
Copy Markdown
Contributor Author

@raulcd Can you run the CI again please?

@Alex-PLACET

Copy link
Copy Markdown
Contributor Author

@raulcd Can you launch the CI please

@Alex-PLACET

Copy link
Copy Markdown
Contributor Author

@raulcd Can you launch the CI please (I hope it will be the last time)

@raulcd

raulcd commented Jul 22, 2026

Copy link
Copy Markdown
Member

@raulcd Can you launch the CI please (I hope it will be the last time)

Sorry @Alex-PLACET , I missed this comment. Just kicked it

@Alex-PLACET
Alex-PLACET force-pushed the refactor/internal-transport branch 2 times, most recently from 560bd8d to 20c2325 Compare July 23, 2026 11:41
Comment on lines +226 to +227
Protobuf_SOURCE=BUNDLED \
gRPC_SOURCE=BUNDLED \

@raulcd raulcd Jul 31, 2026 •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

We should probably bump ARROW_GRPC_REQUIRED_VERSION. I am pretty sure if Ubuntu 22.04 ships older GRPC than necessary, we will also have to update the Linux Package jobs (.deb) for old ubuntu (potentially also for old Red Hat?), I'll kick off Linux Packages.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

maybe the label doesn't work well with draft, PRs, I am unsure why the Linux Package jobs did not trigger, I'll take a look.

::grpc::ByteBuffer* buffer, arrow::flight::internal::FlightData* out);

ARROW_FLIGHT_EXPORT
bool IsRegisteredGrpcFlightDataMessage(

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

we should add some comments on why now we have to register/unregister FlightData messages. Are those necessary?

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.

The point was I got some issue with protobuf and I though it was because of a mess with with FlightPayload. I think I introduce an issue at some point it fixed it later in my development. I just removed this mechanism as it is useless now.

@raulcd raulcd added the CI: Extra: Package: Linux Run extra Linux Packages CI label Jul 31, 2026
@github-actions github-actions Bot added awaiting changes Awaiting changes and removed awaiting review Awaiting review labels Jul 31, 2026

@Alex-PLACET Alex-PLACET Sep 2, 2026 •

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.

Most of the code was moved from transport_server.cc

std::vector<std::shared_ptr<ServerMiddleware>> middleware_;
std::unordered_map<std::string, std::shared_ptr<ServerMiddleware>> middleware_map_;
CallHeaders incoming_headers_;
};

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.

@github-actions github-actions Bot added awaiting change review Awaiting change review and removed awaiting changes Awaiting changes labels Sep 2, 2026
Comment on lines -586 to -621
// Allow uploading messages of any length
builder.SetMaxReceiveMessageSize(-1);

const std::string scheme = uri.scheme();
int port = 0;
if (scheme == kSchemeGrpc || scheme == kSchemeGrpcTcp || scheme == kSchemeGrpcTls) {
std::stringstream address;
address << arrow::util::UriEncodeHost(uri.host()) << ':' << uri.port_text();

std::shared_ptr<::grpc::ServerCredentials> creds;
if (scheme == kSchemeGrpcTls) {
::grpc::SslServerCredentialsOptions ssl_options;
for (const auto& pair : options.tls_certificates) {
ssl_options.pem_key_cert_pairs.push_back({pair.pem_key, pair.pem_cert});
}
if (options.verify_client) {
ssl_options.client_certificate_request =
GRPC_SSL_REQUEST_AND_REQUIRE_CLIENT_CERTIFICATE_AND_VERIFY;
}
if (!options.root_certificates.empty()) {
ssl_options.pem_root_certs = options.root_certificates;
}
creds = ::grpc::SslServerCredentials(ssl_options);
} else {
creds = ::grpc::InsecureServerCredentials();
}

builder.AddListeningPort(address.str(), creds, &port);
} else if (scheme == kSchemeGrpcUnix) {
std::stringstream address;
address << "unix:" << uri.path();
builder.AddListeningPort(address.str(), ::grpc::InsecureServerCredentials());
location_ = options.location;
} else {
return Status::NotImplemented("Scheme is not supported: " + scheme);
}

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.

Moved to a dedicated function: AddServerListeningPort

Comment on lines +38 to +74
Status AddServerListeningPort(const FlightServerOptions& options,
const arrow::util::Uri& uri, ::grpc::ServerBuilder* builder,
Location* location, int* port) {
const std::string scheme = uri.scheme();
if (scheme == kSchemeGrpc || scheme == kSchemeGrpcTcp || scheme == kSchemeGrpcTls) {
std::stringstream address;
address << arrow::util::UriEncodeHost(uri.host()) << ':' << uri.port_text();

std::shared_ptr<::grpc::ServerCredentials> creds;
if (scheme == kSchemeGrpcTls) {
::grpc::SslServerCredentialsOptions ssl_options;
for (const auto& pair : options.tls_certificates) {
ssl_options.pem_key_cert_pairs.push_back({pair.pem_key, pair.pem_cert});
}
if (options.verify_client) {
ssl_options.client_certificate_request =
GRPC_SSL_REQUEST_AND_REQUIRE_CLIENT_CERTIFICATE_AND_VERIFY;
}
if (!options.root_certificates.empty()) {
ssl_options.pem_root_certs = options.root_certificates;
}
creds = ::grpc::SslServerCredentials(ssl_options);
} else {
creds = ::grpc::InsecureServerCredentials();
}
builder->AddListeningPort(address.str(), creds, port);
return Status::OK();
}
if (scheme == kSchemeGrpcUnix) {
std::stringstream address;
address << "unix:" << uri.path();
builder->AddListeningPort(address.str(), ::grpc::InsecureServerCredentials());
*location = options.location;
return Status::OK();
}
return Status::NotImplemented("Scheme is not supported: " + scheme);
}

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.

Moved from grpc_server.cc

}

namespace {
class TransportIpcMessageReader : public ipc::MessageReader {

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.

Moved to transport_server_internal.cc

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.

Most of the code comes from server.cc

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.

Most of the code moved to grpc_server_internal cpp and h

@Alex-PLACET
Alex-PLACET marked this pull request as ready for review September 17, 2026 14:40
Copilot AI lite review requested due to automatic review settings September 17, 2026 14:40
@raulcd raulcd added the CI: Extra: R Run extra R CI label Sep 18, 2026
Copilot AI review requested due to automatic review settings September 18, 2026 08:37

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.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

@Alex-PLACET
Alex-PLACET requested a review from raulcd October 6, 2026 11:27
@pitrou

pitrou commented Oct 6, 2026

Copy link
Copy Markdown
Member

@lidavidm Would you like to take a look at this?

@lidavidm lidavidm left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

It generally seems fine to me.

Comment on lines +201 to +203
ARROW_FLIGHT_EXPORT
Status SetServerLocationFromUri(const arrow::util::Uri& uri, int port,
Location* location);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

nit: this should return Result; it's also not really clear to me why it's factored out


namespace arrow::flight::transport::grpc {

Status AddServerListeningPort(const FlightServerOptions& options,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

nit: it might make a little more sense for this to parse a URI and return the creds/location/address instead of just taking a builder, in case we want to unit test this later

/// Provides shared helpers for constructing message readers/writers from a
/// transport-level data stream, and for writing a FlightDataStream to a
/// transport-level stream (used by DoGet/DoExchange).
class ARROW_FLIGHT_EXPORT ServerTransportBase {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Does this really need to be a base class? It seems WriteDataStream could just be a helper function

Comment on lines +301 to +302
ARROW_FLIGHT_EXPORT
int PortFromLocation(const Location& location);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Could this just be a method on Location?

Comment on lines +293 to +294
ARROW_FLIGHT_EXPORT
arrow::Result<arrow::util::Uri> ParseLocationUri(const Location& location);

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Possibly a method on Location? That would also help resolve the name, which I found a bit confusing

@github-actions github-actions Bot added awaiting review Awaiting review awaiting changes Awaiting changes and removed awaiting change review Awaiting change review labels Oct 7, 2026
- Add Location::uri() and Location::port() and update FlightServerBase::port() to return arrow::Result<int>.
- Remove ParseLocationUri/PortFromLocation and use Location methods instead.
- Remove ServerTransportBase; store MemoryManager on ServerTransport and adjust ctor.
- Make WriteDataStream a free function (transport_server_internal).
- Introduce GrpcServerEndpoint and ParseServerEndpoint (address + credentials + location) and update gRPC server startup to use it.
- Update related headers and implementations to match the new internal API.
Copilot AI lite review requested due to automatic review settings October 7, 2026 14:08
@github-actions github-actions Bot added awaiting review Awaiting review and removed awaiting review Awaiting review awaiting changes Awaiting changes labels Oct 7, 2026

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.

Copilot review overview

🟡 Changes recommended

Unresolved API, compilation, header-installation, gRPC macro, and signal-initialization issues block approval.

Review effort: Lite
Findings: 6 High severity

Open (6)

Comment thread cpp/src/arrow/flight/server.h Outdated
Comment on lines +38 to +42
#define GRPC_CPP_VERSION_CHECK(major, minor, patch) \
((GRPC_CPP_VERSION_MAJOR > (major) || \
(GRPC_CPP_VERSION_MAJOR == (major) && GRPC_CPP_VERSION_MINOR > (minor)) || \
((GRPC_CPP_VERSION_MAJOR == (major) && GRPC_CPP_VERSION_MINOR == (minor) && \
GRPC_CPP_VERSION_PATCH >= (patch)))))
#include <memory>

#include "arrow/flight/transport.h"
#include "arrow/flight/transport_server_internal.h"
Comment on lines +27 to +34
#include "arrow/flight/platform.h"
#include "arrow/flight/server.h"
#include "arrow/flight/visibility.h"
#include "arrow/result.h"
#include "arrow/status.h"
#include "arrow/util/io_util.h"
#include "arrow/util/logging.h"
#include "arrow/util/uri.h"
Comment on lines +202 to +205
shutdown_ = [shutdown = std::forward<ShutdownFn>(shutdown)]() mutable {
return shutdown(nullptr);
};
shutdown_warning_ = shutdown_warning;
/// \brief Get the URI representation of this location.
///
/// \return Arrow result with the URI.
arrow::Result<arrow::util::Uri> uri() const;
Copilot AI lite review requested due to automatic review settings October 7, 2026 14:48

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 build, API compatibility, linking, and signal-handling issues remain.

11 open findings
Previously missed (1)

In code that hasn't changed since last review

Medium severity Return an error when URI port is unavailable

cpp/​src/​arrow/​flight/​types.cc:849

Uri::port() returns -1 when the URI has no port, so this implementation returns a successful Result containing -1 for locations such as grpc+unix. That contradicts the new FlightServerBase::port() contract, which promises an error when the port cannot be determined; convert the no-port sentinel to an error here.

🧠 Review effort: Lite

'server_tracing_middleware.cc',
'transport.cc',
'transport_server.cc',
'transport_server_internal.cc',

#include <sstream>

#include <grpcpp/support/server_callback.h>
Comment on lines +100 to +102
ARROW_FLIGHT_EXPORT
::grpc::Status ToGrpcStatus(const Status& arrow_status,
::grpc::CallbackServerContext* ctx);
#include <memory>

#include "arrow/flight/transport.h"
#include "arrow/flight/transport_server_internal.h"
Comment on lines +213 to +215
running_instance_ = nullptr;
shutdown_ = nullptr;
shutdown_warning_ = nullptr;
Copilot AI lite review requested due to automatic review settings October 7, 2026 15:21

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.

Comment on lines 85 to +91
Status FlightServerBase::Serve() {
if (!impl_->transport_) {
return Status::UnknownError("Server did not start properly");
}
impl_->got_signal_ = 0;
impl_->old_signal_handlers_.clear();
impl_->running_instance_ = impl_.get();

ServerSignalHandler signal_handler;
ARROW_ASSIGN_OR_RAISE(impl_->self_pipe_, signal_handler.Init(&Impl::WaitForSignals));
// Override existing signal handlers with our own handler so as to stop the server.
for (size_t i = 0; i < impl_->signals_.size(); ++i) {
int signum = impl_->signals_[i];
SignalHandler new_handler(&Impl::HandleSignal), old_handler;
ARROW_ASSIGN_OR_RAISE(old_handler, SetSignalHandler(signum, new_handler));
impl_->old_signal_handlers_.push_back(std::move(old_handler));
}

RETURN_NOT_OK(impl_->transport_->Wait());
impl_->running_instance_ = nullptr;

// Restore signal handlers
for (size_t i = 0; i < impl_->signals_.size(); ++i) {
RETURN_NOT_OK(
SetSignalHandler(impl_->signals_[i], impl_->old_signal_handlers_[i]).status());
}
return Status::OK();
return impl_->signal_state_.Serve(
[this]() -> Status {
if (!impl_->transport_) {
return Status::UnknownError("Server did not start properly");
}
return impl_->transport_->Wait();
This reverts commit bcf31a4.
Copilot AI lite review requested due to automatic review settings October 8, 2026 07:06

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.

return impl_->transport_->Shutdown(*deadline);
}
return impl_->transport_->Shutdown();
return impl_->signal_state_.Shutdown(
Comment on lines +38 to +42
#define GRPC_CPP_VERSION_CHECK(major, minor, patch) \
((GRPC_CPP_VERSION_MAJOR > (major) || \
(GRPC_CPP_VERSION_MAJOR == (major) && GRPC_CPP_VERSION_MINOR > (minor)) || \
((GRPC_CPP_VERSION_MAJOR == (major) && GRPC_CPP_VERSION_MINOR == (minor) && \
GRPC_CPP_VERSION_PATCH >= (patch)))))
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants