Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions gen/interfaces/IMapProcessingPipeline_grpcProxy.h
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ class IMapProcessingPipeline_grpcProxy: public org::bcom::xpcf::ConfigurableBas
SolAR::FrameworkReturnCode start() override;
SolAR::FrameworkReturnCode stop() override;
SolAR::FrameworkReturnCode setMapToProcess(SRef<SolAR::datastructure::Map> const map) override;
SolAR::FrameworkReturnCode setMapToProcess(std::string const& mapUUID, std::string const& resultMapUUID) override;
SolAR::FrameworkReturnCode getStatus(SolAR::api::pipeline::MapProcessingStatus& status, float& progress) const override;
SolAR::FrameworkReturnCode getProcessingData(std::vector<SRef<SolAR::datastructure::CloudPoint>>& pointCloud, std::vector<SolAR::datastructure::Transform3Df>& keyframePoses) const override;
SolAR::FrameworkReturnCode getProcessedMap(SRef<SolAR::datastructure::Map>& map) const override;
Expand Down
3 changes: 2 additions & 1 deletion gen/interfaces/IMapProcessingPipeline_grpcServer.h
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,8 @@ class IMapProcessingPipeline_grpcServer: public org::bcom::xpcf::ConfigurableBa
::grpc::Status init(::grpc::ServerContext* context, const ::grpcIMapProcessingPipeline::initRequest* request, ::grpcIMapProcessingPipeline::initResponse* response) override;
::grpc::Status start(::grpc::ServerContext* context, const ::grpcIMapProcessingPipeline::startRequest* request, ::grpcIMapProcessingPipeline::startResponse* response) override;
::grpc::Status stop(::grpc::ServerContext* context, const ::grpcIMapProcessingPipeline::stopRequest* request, ::grpcIMapProcessingPipeline::stopResponse* response) override;
::grpc::Status setMapToProcess(::grpc::ServerContext* context, const ::grpcIMapProcessingPipeline::setMapToProcessRequest* request, ::grpcIMapProcessingPipeline::setMapToProcessResponse* response) override;
::grpc::Status setMapToProcess_grpc0(::grpc::ServerContext* context, const ::grpcIMapProcessingPipeline::setMapToProcess_grpc0Request* request, ::grpcIMapProcessingPipeline::setMapToProcess_grpc0Response* response) override;
::grpc::Status setMapToProcess_grpc1(::grpc::ServerContext* context, const ::grpcIMapProcessingPipeline::setMapToProcess_grpc1Request* request, ::grpcIMapProcessingPipeline::setMapToProcess_grpc1Response* response) override;
::grpc::Status getStatus(::grpc::ServerContext* context, const ::grpcIMapProcessingPipeline::getStatusRequest* request, ::grpcIMapProcessingPipeline::getStatusResponse* response) override;
::grpc::Status getProcessingData(::grpc::ServerContext* context, const ::grpcIMapProcessingPipeline::getProcessingDataRequest* request, ::grpcIMapProcessingPipeline::getProcessingDataResponse* response) override;
::grpc::Status getProcessedMap(::grpc::ServerContext* context, const ::grpcIMapProcessingPipeline::getProcessedMapRequest* request, ::grpcIMapProcessingPipeline::getProcessedMapResponse* response) override;
Expand Down
337 changes: 247 additions & 90 deletions gen/interfaces/grpcIMapProcessingPipelineService.grpc.pb.h

Large diffs are not rendered by default.

826 changes: 718 additions & 108 deletions gen/interfaces/grpcIMapProcessingPipelineService.pb.h

Large diffs are not rendered by default.

19 changes: 16 additions & 3 deletions gen/proto/grpcIMapProcessingPipelineService.proto
Original file line number Diff line number Diff line change
Expand Up @@ -34,13 +34,25 @@ message stopResponse
sint32 xpcfGrpcReturnValue = 1;
}

message setMapToProcessRequest
message setMapToProcess_grpc0Request
{
int32 grpcServerCompressionFormat = 1;
bytes map = 2;
}

message setMapToProcessResponse
message setMapToProcess_grpc0Response
{
sint32 xpcfGrpcReturnValue = 1;
}

message setMapToProcess_grpc1Request
{
int32 grpcServerCompressionFormat = 1;
string mapUUID = 2;
string resultMapUUID = 3;
}

message setMapToProcess_grpc1Response
{
sint32 xpcfGrpcReturnValue = 1;
}
Expand Down Expand Up @@ -89,7 +101,8 @@ service grpcIMapProcessingPipelineService {
rpc init(initRequest) returns(initResponse) {}
rpc start(startRequest) returns(startResponse) {}
rpc stop(stopRequest) returns(stopResponse) {}
rpc setMapToProcess(setMapToProcessRequest) returns(setMapToProcessResponse) {}
rpc setMapToProcess_grpc0(setMapToProcess_grpc0Request) returns(setMapToProcess_grpc0Response) {}
rpc setMapToProcess_grpc1(setMapToProcess_grpc1Request) returns(setMapToProcess_grpc1Response) {}
rpc getStatus(getStatusRequest) returns(getStatusResponse) {}
rpc getProcessingData(getProcessingDataRequest) returns(getProcessingDataResponse) {}
rpc getProcessedMap(getProcessedMapRequest) returns(getProcessedMapResponse) {}
Expand Down
80 changes: 72 additions & 8 deletions gen/src/IMapProcessingPipeline_grpcProxy.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ IMapProcessingPipeline_grpcProxy::IMapProcessingPipeline_grpcProxy():xpcf::Confi
declareInterface<SolAR::api::pipeline::IMapProcessingPipeline>(this);
declareProperty("channelUrl",m_channelUrl);
declareProperty("channelCredentials",m_channelCredentials);
m_grpcProxyCompressionConfig.resize(8);
m_grpcProxyCompressionConfig.resize(9);
declarePropertySequence("grpc_compress_proxy", m_grpcProxyCompressionConfig);
}

Expand Down Expand Up @@ -270,8 +270,8 @@ SolAR::FrameworkReturnCode IMapProcessingPipeline_grpcProxy::stop()
SolAR::FrameworkReturnCode IMapProcessingPipeline_grpcProxy::setMapToProcess(SRef<SolAR::datastructure::Map> const map)
{
::grpc::ClientContext context;
::grpcIMapProcessingPipeline::setMapToProcessRequest reqIn;
::grpcIMapProcessingPipeline::setMapToProcessResponse respOut;
::grpcIMapProcessingPipeline::setMapToProcess_grpc0Request reqIn;
::grpcIMapProcessingPipeline::setMapToProcess_grpc0Response respOut;
#ifndef DISABLE_GRPC_COMPRESSION
xpcf::grpcCompressionInfos proxyCompressionInfo = xpcf::deduceClientCompressionInfo(m_serviceCompressionInfos, "setMapToProcess", m_methodCompressionInfosMap);
xpcf::grpcCompressType serverCompressionType = xpcf::prepareClientCompressionContext(context, proxyCompressionInfo);
Expand All @@ -296,7 +296,7 @@ SolAR::FrameworkReturnCode IMapProcessingPipeline_grpcProxy::setMapToProcess(SR
auto span = tracer->StartSpan("IMapProcessingPipeline_grpcProxy.setMapToProcess",
{{opentelemetry::semconv::rpc::kRpcSystem, "grpc"},
{opentelemetry::semconv::rpc::kRpcService, "grpcIMapProcessingPipeline.grpcIMapProcessingPipelineService"},
{opentelemetry::semconv::rpc::kRpcMethod, "setMapToProcess"},
{opentelemetry::semconv::rpc::kRpcMethod, "setMapToProcess_grpc0"},
{opentelemetry::semconv::network::kNetworkPeerAddress, networkAddress},
{opentelemetry::semconv::network::kNetworkPeerPort, std::stoi(networkPort)}},
spanOptions);
Expand All @@ -309,17 +309,81 @@ SolAR::FrameworkReturnCode IMapProcessingPipeline_grpcProxy::setMapToProcess(SR
auto prop = opentelemetry::context::propagation::GlobalTextMapPropagator::GetGlobalPropagator();
prop->Inject(carrier, currentCtx);

::grpc::Status grpcRemoteStatus = m_grpcStub->setMapToProcess(&context, reqIn, &respOut);
::grpc::Status grpcRemoteStatus = m_grpcStub->setMapToProcess_grpc0(&context, reqIn, &respOut);
#ifdef ENABLE_PROXY_TIMERS
boost::posix_time::ptime end = boost::posix_time::microsec_clock::universal_time();
std::cout << "====> IMapProcessingPipeline_grpcProxy::setMapToProcess response received at " << to_simple_string(end) << std::endl;
std::cout << " => elapsed time = " << ((end - start).total_microseconds() / 1000.00) << " ms" << std::endl;
#endif
if (!grpcRemoteStatus.ok()) {
std::cout << "setMapToProcess rpc failed." << std::endl;
span->SetStatus(opentelemetry::trace::StatusCode::kError, "grpcIMapProcessingPipelineService.setMapToProcess() rpc failed.");
std::cout << "setMapToProcess_grpc0 rpc failed." << std::endl;
span->SetStatus(opentelemetry::trace::StatusCode::kError, "grpcIMapProcessingPipelineService.setMapToProcess_grpc0() rpc failed.");
span->End();
throw xpcf::RemotingException("grpcIMapProcessingPipelineService","setMapToProcess",static_cast<uint32_t>(grpcRemoteStatus.error_code()));
throw xpcf::RemotingException("grpcIMapProcessingPipelineService","setMapToProcess_grpc0",static_cast<uint32_t>(grpcRemoteStatus.error_code()));
}


span->SetStatus(opentelemetry::trace::StatusCode::kOk);
span->End();

return static_cast<SolAR::FrameworkReturnCode>(respOut.xpcfgrpcreturnvalue());
}


SolAR::FrameworkReturnCode IMapProcessingPipeline_grpcProxy::setMapToProcess(std::string const& mapUUID, std::string const& resultMapUUID)
{
::grpc::ClientContext context;
::grpcIMapProcessingPipeline::setMapToProcess_grpc1Request reqIn;
::grpcIMapProcessingPipeline::setMapToProcess_grpc1Response respOut;
#ifndef DISABLE_GRPC_COMPRESSION
xpcf::grpcCompressionInfos proxyCompressionInfo = xpcf::deduceClientCompressionInfo(m_serviceCompressionInfos, "setMapToProcess", m_methodCompressionInfosMap);
xpcf::grpcCompressType serverCompressionType = xpcf::prepareClientCompressionContext(context, proxyCompressionInfo);
reqIn.set_grpcservercompressionformat (static_cast<int32_t>(serverCompressionType));
#endif
reqIn.set_mapuuid(mapUUID);
reqIn.set_resultmapuuid(resultMapUUID);
#ifdef ENABLE_PROXY_TIMERS
boost::posix_time::ptime start = boost::posix_time::microsec_clock::universal_time();
std::cout << "====> IMapProcessingPipeline_grpcProxy::setMapToProcess request sent at " << to_simple_string(start) << std::endl;
#endif

auto provider = opentelemetry::trace::Provider::GetTracerProvider();
auto tracer = provider->GetTracer("xpcfGrpcRemotingSolARFramework", "1.6.0");

// TODO: safer parsing with error handling
auto const pos = m_channelUrl.find_last_of(':');
auto networkAddress = m_channelUrl.substr(0, pos);
auto networkPort = m_channelUrl.substr(pos + 1);

opentelemetry::trace::StartSpanOptions spanOptions;
spanOptions.kind = opentelemetry::trace::SpanKind::kClient;
auto span = tracer->StartSpan("IMapProcessingPipeline_grpcProxy.setMapToProcess",
{{opentelemetry::semconv::rpc::kRpcSystem, "grpc"},
{opentelemetry::semconv::rpc::kRpcService, "grpcIMapProcessingPipeline.grpcIMapProcessingPipelineService"},
{opentelemetry::semconv::rpc::kRpcMethod, "setMapToProcess_grpc1"},
{opentelemetry::semconv::network::kNetworkPeerAddress, networkAddress},
{opentelemetry::semconv::network::kNetworkPeerPort, std::stoi(networkPort)}},
spanOptions);

auto scope = tracer->WithActiveSpan(span);

// inject current context to grpc metadata
auto currentCtx = opentelemetry::context::RuntimeContext::GetCurrent();
GrpcClientCarrier carrier(&context);
auto prop = opentelemetry::context::propagation::GlobalTextMapPropagator::GetGlobalPropagator();
prop->Inject(carrier, currentCtx);

::grpc::Status grpcRemoteStatus = m_grpcStub->setMapToProcess_grpc1(&context, reqIn, &respOut);
#ifdef ENABLE_PROXY_TIMERS
boost::posix_time::ptime end = boost::posix_time::microsec_clock::universal_time();
std::cout << "====> IMapProcessingPipeline_grpcProxy::setMapToProcess response received at " << to_simple_string(end) << std::endl;
std::cout << " => elapsed time = " << ((end - start).total_microseconds() / 1000.00) << " ms" << std::endl;
#endif
if (!grpcRemoteStatus.ok()) {
std::cout << "setMapToProcess_grpc1 rpc failed." << std::endl;
span->SetStatus(opentelemetry::trace::StatusCode::kError, "grpcIMapProcessingPipelineService.setMapToProcess_grpc1() rpc failed.");
span->End();
throw xpcf::RemotingException("grpcIMapProcessingPipelineService","setMapToProcess_grpc1",static_cast<uint32_t>(grpcRemoteStatus.error_code()));
}


Expand Down
51 changes: 48 additions & 3 deletions gen/src/IMapProcessingPipeline_grpcServer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -95,7 +95,7 @@ IMapProcessingPipeline_grpcServer::IMapProcessingPipeline_grpcServer():xpcf::Con
{
declareInterface<xpcf::IGrpcService>(this);
declareInjectable<SolAR::api::pipeline::IMapProcessingPipeline>(m_grpcService.m_xpcfComponent);
m_grpcServerCompressionConfig.resize(8);
m_grpcServerCompressionConfig.resize(9);
declarePropertySequence("grpc_compress_server", m_grpcServerCompressionConfig);
}

Expand Down Expand Up @@ -250,7 +250,7 @@ ::grpc::Status IMapProcessingPipeline_grpcServer::grpcIMapProcessingPipelineServ
}


::grpc::Status IMapProcessingPipeline_grpcServer::grpcIMapProcessingPipelineServiceImpl::setMapToProcess(::grpc::ServerContext* context, const ::grpcIMapProcessingPipeline::setMapToProcessRequest* request, ::grpcIMapProcessingPipeline::setMapToProcessResponse* response)
::grpc::Status IMapProcessingPipeline_grpcServer::grpcIMapProcessingPipelineServiceImpl::setMapToProcess_grpc0(::grpc::ServerContext* context, const ::grpcIMapProcessingPipeline::setMapToProcess_grpc0Request* request, ::grpcIMapProcessingPipeline::setMapToProcess_grpc0Response* response)
{
auto prop = opentelemetry::context::propagation::GlobalTextMapPropagator::GetGlobalPropagator();
auto currentCtx = opentelemetry::context::RuntimeContext::GetCurrent();
Expand All @@ -267,7 +267,7 @@ ::grpc::Status IMapProcessingPipeline_grpcServer::grpcIMapProcessingPipelineServ
auto span = tracer->StartSpan("IMapProcessingPipeline_grpcServer.setMapToProcess",
{{opentelemetry::semconv::rpc::kRpcSystem, "grpc"},
{opentelemetry::semconv::rpc::kRpcService, "grpcIMapProcessingPipeline.grpcIMapProcessingPipelineService"},
{opentelemetry::semconv::rpc::kRpcMethod, "setMapToProcess"},
{opentelemetry::semconv::rpc::kRpcMethod, "setMapToProcess_grpc0"},
{opentelemetry::semconv::rpc::kRpcGrpcStatusCode, 0}},
options);
SpanScope spanScope(span);
Expand All @@ -294,6 +294,51 @@ ::grpc::Status IMapProcessingPipeline_grpcServer::grpcIMapProcessingPipelineServ
}


::grpc::Status IMapProcessingPipeline_grpcServer::grpcIMapProcessingPipelineServiceImpl::setMapToProcess_grpc1(::grpc::ServerContext* context, const ::grpcIMapProcessingPipeline::setMapToProcess_grpc1Request* request, ::grpcIMapProcessingPipeline::setMapToProcess_grpc1Response* response)
{
auto prop = opentelemetry::context::propagation::GlobalTextMapPropagator::GetGlobalPropagator();
auto currentCtx = opentelemetry::context::RuntimeContext::GetCurrent();
GrpcServerCarrier carrier(context);
auto newContext = prop->Extract(carrier, currentCtx);
ContextScope ctxtScope(newContext);

opentelemetry::trace::StartSpanOptions options;
options.kind = opentelemetry::trace::SpanKind::kServer;
options.parent = opentelemetry::trace::GetSpan(newContext)->GetContext();

auto provider = opentelemetry::trace::Provider::GetTracerProvider();
auto tracer = provider->GetTracer("xpcfGrpcRemotingSolARFramework", "1.6.0");
auto span = tracer->StartSpan("IMapProcessingPipeline_grpcServer.setMapToProcess",
{{opentelemetry::semconv::rpc::kRpcSystem, "grpc"},
{opentelemetry::semconv::rpc::kRpcService, "grpcIMapProcessingPipeline.grpcIMapProcessingPipelineService"},
{opentelemetry::semconv::rpc::kRpcMethod, "setMapToProcess_grpc1"},
{opentelemetry::semconv::rpc::kRpcGrpcStatusCode, 0}},
options);
SpanScope spanScope(span);
auto scope= tracer->WithActiveSpan(span);

#ifndef DISABLE_GRPC_COMPRESSION
xpcf::grpcCompressType askedCompressionType = static_cast<xpcf::grpcCompressType>(request->grpcservercompressionformat());
xpcf::grpcServerCompressionInfos serverCompressInfo = xpcf::deduceServerCompressionType(askedCompressionType, m_serviceCompressionInfos, "setMapToProcess", m_methodCompressionInfosMap);
xpcf::prepareServerCompressionContext(context, serverCompressInfo);
#endif
#ifdef ENABLE_SERVER_TIMERS
boost::posix_time::ptime start = boost::posix_time::microsec_clock::universal_time();
std::cout << "====> IMapProcessingPipeline_grpcServer::setMapToProcess request received at " << to_simple_string(start) << std::endl;
#endif
std::string mapUUID = request->mapuuid();
std::string resultMapUUID = request->resultmapuuid();
SolAR::FrameworkReturnCode returnValue = m_xpcfComponent->setMapToProcess(mapUUID, resultMapUUID);
response->set_xpcfgrpcreturnvalue(static_cast<int32_t>(returnValue));
#ifdef ENABLE_SERVER_TIMERS
boost::posix_time::ptime end = boost::posix_time::microsec_clock::universal_time();
std::cout << "====> IMapProcessingPipeline_grpcServer::setMapToProcess response sent at " << to_simple_string(end) << std::endl;
std::cout << " => elapsed time = " << ((end - start).total_microseconds() / 1000.00) << " ms" << std::endl;
#endif
return ::grpc::Status::OK;
}


::grpc::Status IMapProcessingPipeline_grpcServer::grpcIMapProcessingPipelineServiceImpl::getStatus(::grpc::ServerContext* context, const ::grpcIMapProcessingPipeline::getStatusRequest* request, ::grpcIMapProcessingPipeline::getStatusResponse* response)
{
auto prop = opentelemetry::context::propagation::GlobalTextMapPropagator::GetGlobalPropagator();
Expand Down
Loading