#include "wilson_uploader.h" #include #include #include #include #include #include #include namespace NWilson { using namespace NActors; namespace NServiceProto = opentelemetry::proto::collector::trace::v1; namespace { class TWilsonUploader : public TActorBootstrapped { TString Host; ui16 Port; TString RootCA; std::shared_ptr Channel; std::unique_ptr Stub; grpc::CompletionQueue CQ; std::unique_ptr Context; std::unique_ptr> Reader; NServiceProto::ExportTraceServiceRequest Request; NServiceProto::ExportTraceServiceResponse Response; grpc::Status Status; public: TWilsonUploader(TString host, ui16 port, TString rootCA) : Host(std::move(host)) , Port(std::move(port)) , RootCA(std::move(rootCA)) {} ~TWilsonUploader() { CQ.Shutdown(); } void Bootstrap() { Become(&TThis::StateFunc); Channel = grpc::CreateChannel(TStringBuilder() << Host << ":" << Port, grpc::SslCredentials({ .pem_root_certs = TFileInput(RootCA).ReadAll(), })); Stub = NServiceProto::TraceService::NewStub(Channel); } void Handle(TEvWilson::TPtr ev) { CheckIfDone(); auto *rspan = Request.resource_spans_size() ? Request.mutable_resource_spans(0) : Request.add_resource_spans(); auto *sspan = rspan->scope_spans_size() ? rspan->mutable_scope_spans(0) : rspan->add_scope_spans(); ev->Get()->Span.Swap(sspan->add_spans()); if (!Context) { SendRequest(); } } void SendRequest() { Y_VERIFY(!Reader && !Context); Context = std::make_unique(); Reader = Stub->AsyncExport(Context.get(), std::exchange(Request, {}), &CQ); Reader->Finish(&Response, &Status, nullptr); } void CheckIfDone() { if (Context) { void *tag; bool ok; if (CQ.AsyncNext(&tag, &ok, std::chrono::system_clock::now()) == grpc::CompletionQueue::GOT_EVENT) { Reader.reset(); Context.reset(); if (Request.resource_spans_size()) { SendRequest(); } } } } STRICT_STFUNC(StateFunc, hFunc(TEvWilson, Handle); ); }; } // anonymous IActor *CreateWilsonUploader(TString host, ui16 port, TString rootCA) { return new TWilsonUploader(std::move(host), port, std::move(rootCA)); } } // NWilson