Skip to content

Commit 98c413a

Browse files
committed
agents: add async OTLP GrpcMetricsExporter
Introduce GrpcMetricsExporter that uses DelegateAsyncExport + AsyncTSQueue to flush MetricDataBatch instances off the main thread and retry transient gRPC failures with backpressure. Add NSOLID_METRICS_BUFFER_SIZE option (plumbed through lib/nsolid.js, reconfigure proto/pb.cc/pb.h, and node.gyp) plus new test-grpc-metrics-buffer-size.mjs and retry coverage. It defaults to 100 elements, following the ZmqAgent behavior. Wire GrpcAgent config/reconfigure paths to construct the exporter once and call resize_buffer on buffer-size changes, ensuring the new queue depth propagates without losing metrics.
1 parent 2c31e0b commit 98c413a

14 files changed

Lines changed: 592 additions & 57 deletions

agents/grpc/proto/reconfigure.proto

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@ message ReconfigureBody {
1919
optional bool contCpuProfile = 12;
2020
optional bool assetsEnabled = 13;
2121
optional uint32 metricsBatchSize = 14;
22+
optional uint32 metricsBufferSize = 15;
2223
}
2324

2425
message ReconfigureEvent {

agents/grpc/src/grpc_agent.cc

Lines changed: 34 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,6 @@
88
#include "../../src/span_collector.h"
99
#include "absl/log/initialize.h"
1010
#include "opentelemetry/sdk/metrics/data/metric_data.h"
11-
#include "opentelemetry/sdk/metrics/export/metric_producer.h"
1211
#include "opentelemetry/semconv/incubating/process_attributes.h"
1312
#include "opentelemetry/semconv/service_attributes.h"
1413
#include "opentelemetry/exporters/otlp/otlp_grpc_client.h"
@@ -17,8 +16,6 @@
1716
#include "opentelemetry/exporters/otlp/otlp_grpc_exporter_factory.h"
1817
#include "opentelemetry/exporters/otlp/otlp_grpc_log_record_exporter.h"
1918
#include "opentelemetry/exporters/otlp/otlp_grpc_log_record_exporter_factory.h"
20-
#include "opentelemetry/exporters/otlp/otlp_grpc_metric_exporter.h"
21-
#include "opentelemetry/exporters/otlp/otlp_grpc_metric_exporter_factory.h"
2219
#include "opentelemetry/exporters/otlp/otlp_metric_utils.h"
2320

2421
using std::chrono::duration_cast;
@@ -49,8 +46,6 @@ using opentelemetry::v1::exporter::otlp::OtlpGrpcExporterOptions;
4946
using opentelemetry::v1::exporter::otlp::OtlpGrpcLogRecordExporter;
5047
using opentelemetry::v1::exporter::otlp::OtlpGrpcLogRecordExporterFactory;
5148
using opentelemetry::v1::exporter::otlp::OtlpGrpcLogRecordExporterOptions;
52-
using opentelemetry::v1::exporter::otlp::OtlpGrpcMetricExporter;
53-
using opentelemetry::v1::exporter::otlp::OtlpGrpcMetricExporterFactory;
5449
using opentelemetry::v1::exporter::otlp::OtlpGrpcMetricExporterOptions;
5550
using opentelemetry::v1::exporter::otlp::OtlpMetricUtils;
5651
using nsolid_grpc_async =
@@ -955,8 +950,7 @@ void GrpcAgent::env_deletion_cb_(SharedEnvInst envinst,
955950
data.resource_ = otlp::GetResource();
956951
data.scope_metric_data_ =
957952
std::vector<ScopeMetrics>{{ otlp::GetScope(), std::move(metric_data) }};
958-
auto result = agent->metrics_exporter_->Export(data);
959-
Debug("# ThreadMetrics Exported. Result: %d\n", static_cast<int>(result));
953+
agent->metrics_exporter_->enqueue(std::move(data));
960954
}
961955
}
962956
}
@@ -979,9 +973,7 @@ void GrpcAgent::on_metrics_timer() {
979973
data.resource_ = otlp::GetResource();
980974
data.scope_metric_data_ =
981975
std::vector<ScopeMetrics>{{ otlp::GetScope(), std::move(metric_data) }};
982-
auto result = metrics_exporter_->Export(data);
983-
Debug("# ProcessMetrics Exported. Result: %d\n",
984-
static_cast<int>(result));
976+
metrics_exporter_->enqueue(std::move(data));
985977
}
986978

987979
metric_data = thr_metrics_batch_.DumpMetricsAndReset();
@@ -990,10 +982,10 @@ void GrpcAgent::on_metrics_timer() {
990982
data.resource_ = otlp::GetResource();
991983
data.scope_metric_data_ =
992984
std::vector<ScopeMetrics>{{ otlp::GetScope(), std::move(metric_data) }};
993-
auto result = metrics_exporter_->Export(data);
994-
Debug("# ThreadMetrics Exported. Result: %d\n", static_cast<int>(result));
985+
metrics_exporter_->enqueue(std::move(data));
995986
}
996987

988+
metrics_exporter_->flush();
997989
return;
998990
}
999991

@@ -1150,8 +1142,18 @@ int GrpcAgent::config(const json& config) {
11501142
options.credentials = opts.credentials;
11511143
}
11521144

1153-
metrics_exporter_ =
1154-
std::make_unique<OtlpGrpcMetricExporter>(options, client);
1145+
// Get metrics buffer size from config, default to 100
1146+
size_t buffer_size = 100;
1147+
auto it = config_.find("metricsBufferSize");
1148+
if (it != config_.end()) {
1149+
buffer_size = it->get<size_t>();
1150+
}
1151+
1152+
metrics_exporter_ = std::make_unique<GrpcMetricsExporter>(&loop_,
1153+
options,
1154+
client,
1155+
buffer_size);
1156+
metrics_exporter_->init();
11551157
}
11561158
{
11571159
OtlpGrpcLogRecordExporterOptions options;
@@ -1222,6 +1224,17 @@ int GrpcAgent::config(const json& config) {
12221224
}
12231225
}
12241226

1227+
if (utils::find_any_fields_in_diff(diff, { "/metricsBufferSize" })) {
1228+
auto it = config_.find("metricsBufferSize");
1229+
if (it != config_.end()) {
1230+
size_t buffer_size = it->get<size_t>();
1231+
if (buffer_size > 0 && metrics_exporter_) {
1232+
// Resize the metrics buffer instead of recreating the exporter
1233+
metrics_exporter_->resize_buffer(buffer_size);
1234+
}
1235+
}
1236+
}
1237+
12251238
// If metrics timer is not active or if the diff contains metrics fields,
12261239
// recalculate the metrics status. (stop/start/what period)
12271240
if (!metrics_timer_.is_active() ||
@@ -1235,6 +1248,8 @@ int GrpcAgent::config(const json& config) {
12351248
if (it != config_.end()) {
12361249
period = *it;
12371250
}
1251+
} else {
1252+
period = 5000;
12381253
}
12391254
}
12401255

@@ -1356,6 +1371,7 @@ void GrpcAgent::do_stop() {
13561371
}
13571372

13581373
log_exporter_.reset();
1374+
otlp_grpc_client_.reset();
13591375
metrics_exporter_.reset();
13601376
trace_exporter_.reset();
13611377
ready_ = false;
@@ -1447,9 +1463,7 @@ void GrpcAgent::got_proc_metrics() {
14471463
data.resource_ = otlp::GetResource();
14481464
data.scope_metric_data_ =
14491465
std::vector<ScopeMetrics>{{ otlp::GetScope(), std::move(metric_data) }};
1450-
auto result = metrics_exporter_->Export(data);
1451-
Debug("# ProcessMetrics Exported. Result: %d\n",
1452-
static_cast<int>(result));
1466+
metrics_exporter_->enqueue(std::move(data));
14531467
}
14541468
}
14551469
}
@@ -1826,6 +1840,9 @@ void GrpcAgent::reconfigure(const grpcagent::CommandRequest& request) {
18261840
if (body.has_metricsbatchsize()) {
18271841
out["metricsBatchSize"] = body.metricsbatchsize();
18281842
}
1843+
if (body.has_metricsbuffersize()) {
1844+
out["metricsBufferSize"] = body.metricsbuffersize();
1845+
}
18291846

18301847
DebugJSON("Reconfigure out: \n%s\n", out);
18311848

agents/grpc/src/grpc_agent.h

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -9,12 +9,14 @@
99
#include "./proto/nsolid_service.grpc.pb.h"
1010
#include "opentelemetry/version.h"
1111
#include "opentelemetry/sdk/trace/recordable.h"
12+
#include "opentelemetry/exporters/otlp/otlp_grpc_client.h"
1213
#include "../../otlp/src/otlp_common.h"
1314
#include "../../src/profile_collector.h"
1415
#include "asset_stream.h"
1516
#include "command_stream.h"
1617
#include "grpc_client.h"
1718
#include "grpc_errors.h"
19+
#include "grpc_metrics_exporter.h"
1820

1921
// Class pre-declaration
2022
OPENTELEMETRY_BEGIN_NAMESPACE
@@ -148,6 +150,8 @@ class GrpcAgent: public std::enable_shared_from_this<GrpcAgent>,
148150
bool testing;
149151
};
150152

153+
static constexpr uint32_t kMaxMetricsExportRetries = 5;
154+
151155
GrpcAgent();
152156

153157
~GrpcAgent();
@@ -187,6 +191,8 @@ class GrpcAgent: public std::enable_shared_from_this<GrpcAgent>,
187191

188192
static void metrics_timer_cb_(nsuv::ns_timer*, WeakGrpcAgent);
189193

194+
static void metrics_retry_cb_(nsuv::ns_async*, WeakGrpcAgent);
195+
190196
static void profile_msg_cb_(nsuv::ns_async*, WeakGrpcAgent);
191197

192198
static void shutdown_cb_(nsuv::ns_async*, WeakGrpcAgent);
@@ -313,6 +319,11 @@ class GrpcAgent: public std::enable_shared_from_this<GrpcAgent>,
313319
// Blocked Loop
314320
std::shared_ptr<AsyncTSQueue<BlockedLoopStor>> blocked_loop_queue_;
315321

322+
std::shared_ptr<opentelemetry::v1::exporter::otlp::OtlpGrpcClient>
323+
otlp_grpc_client_;
324+
opentelemetry::v1::exporter::otlp::OtlpGrpcClientOptions
325+
otlp_grpc_client_options_;
326+
316327
// For the Tracing API
317328
uint32_t trace_flags_;
318329
std::shared_ptr<SpanCollector> span_collector_;
@@ -323,18 +334,16 @@ class GrpcAgent: public std::enable_shared_from_this<GrpcAgent>,
323334

324335
// For the Metrics API
325336
bool metrics_paused_;
326-
uint64_t metrics_interval_;
327337
ProcessMetrics proc_metrics_;
328338
ProcessMetrics::MetricsStor proc_prev_stor_;
329339
otlp::MetricDataBatch proc_metrics_batch_;
330340
std::map<uint64_t, JSThreadMetrics> env_metrics_map_;
331341
nsuv::ns_async metrics_msg_;
332342
TSQueue<ThreadMetrics::MetricsStor> thr_metrics_msg_q_;
333343
nsuv::ns_timer metrics_timer_;
334-
std::unique_ptr<opentelemetry::v1::exporter::otlp::OtlpGrpcMetricExporter>
335-
metrics_exporter_;
336344
std::map<uint64_t, ThreadMetrics::MetricsStor> thr_metrics_cache_;
337345
otlp::MetricDataBatch thr_metrics_batch_;
346+
std::shared_ptr<GrpcMetricsExporter> metrics_exporter_;
338347

339348
// For the Configuration API
340349
nsuv::ns_async config_msg_;
Lines changed: 128 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,128 @@
1+
#include "grpc_metrics_exporter.h"
2+
#include "grpc_client.h"
3+
4+
#include "opentelemetry/exporters/otlp/otlp_grpc_client.h"
5+
#include "opentelemetry/exporters/otlp/otlp_metric_utils.h"
6+
7+
using
8+
opentelemetry::proto::collector::metrics::v1::ExportMetricsServiceRequest;
9+
using
10+
opentelemetry::proto::collector::metrics::v1::ExportMetricsServiceResponse;
11+
using opentelemetry::proto::collector::metrics::v1::MetricsService;
12+
using opentelemetry::sdk::common::ExportResult;
13+
using opentelemetry::sdk::metrics::ResourceMetrics;
14+
using opentelemetry::v1::exporter::otlp::OtlpGrpcClient;
15+
using opentelemetry::v1::exporter::otlp::OtlpGrpcMetricExporterOptions;
16+
using opentelemetry::v1::exporter::otlp::OtlpMetricUtils;
17+
18+
namespace node {
19+
namespace nsolid {
20+
namespace grpc {
21+
22+
using SharedGrpcMetricsExporter = std::shared_ptr<GrpcMetricsExporter>;
23+
using WeakGrpcMetricsExporter = std::weak_ptr<GrpcMetricsExporter>;
24+
25+
GrpcMetricsExporter::GrpcMetricsExporter(
26+
uv_loop_t* loop,
27+
OtlpGrpcMetricExporterOptions options,
28+
std::shared_ptr<OtlpGrpcClient> client,
29+
size_t buffer_size):
30+
loop_(loop),
31+
options_(options),
32+
client_(client),
33+
metrics_service_stub_(client_->MakeMetricsServiceStub()),
34+
metrics_q_(buffer_size),
35+
in_flight_(false) {
36+
}
37+
38+
void GrpcMetricsExporter::init() {
39+
metrics_completion_q_ = AsyncTSQueue<::grpc::Status>::create(
40+
loop_,
41+
+[](::grpc::Status status, WeakGrpcMetricsExporter exporter_wp) {
42+
SharedGrpcMetricsExporter exporter = exporter_wp.lock();
43+
if (exporter == nullptr) {
44+
return;
45+
}
46+
47+
exporter->on_metrics_export_complete(status);
48+
},
49+
weak_from_this());
50+
}
51+
52+
GrpcMetricsExporter::~GrpcMetricsExporter() {
53+
}
54+
55+
void GrpcMetricsExporter::enqueue(ResourceMetrics&& metrics) {
56+
metrics_q_.push(std::move(metrics));
57+
if (!in_flight_) {
58+
export_current();
59+
}
60+
}
61+
62+
void GrpcMetricsExporter::export_current() {
63+
if (metrics_q_.empty()) {
64+
return;
65+
}
66+
67+
auto* stub = metrics_service_stub_.get();
68+
if (stub == nullptr) {
69+
return;
70+
}
71+
72+
auto& data = metrics_q_.front();
73+
74+
google::protobuf::ArenaOptions arena_options;
75+
arena_options.initial_block_size = 1024;
76+
arena_options.max_block_size = 65536;
77+
auto arena = std::make_unique<google::protobuf::Arena>(arena_options);
78+
79+
auto* request =
80+
google::protobuf::Arena::Create<ExportMetricsServiceRequest>(arena.get());
81+
OtlpMetricUtils::PopulateRequest(data, request);
82+
83+
auto context = OtlpGrpcClient::MakeClientContext(options_);
84+
::grpc::Status immediate =
85+
GrpcClient::DelegateAsyncExport<MetricsServiceStub,
86+
ExportMetricsServiceRequest,
87+
ExportMetricsServiceResponse>(
88+
stub,
89+
&MetricsServiceStub::async_interface::Export,
90+
std::move(context),
91+
std::move(arena),
92+
std::move(*request),
93+
[weak = weak_from_this()](::grpc::Status status,
94+
std::unique_ptr<google::protobuf::Arena>&&,
95+
const ExportMetricsServiceRequest&,
96+
ExportMetricsServiceResponse*) {
97+
auto exporter = weak.lock();
98+
if (exporter != nullptr) {
99+
exporter->metrics_completion_q_->enqueue(status);
100+
}
101+
});
102+
103+
if (!immediate.ok()) {
104+
metrics_completion_q_->enqueue(immediate);
105+
} else {
106+
in_flight_ = true;
107+
}
108+
}
109+
110+
void GrpcMetricsExporter::flush() {
111+
export_current();
112+
}
113+
114+
void GrpcMetricsExporter::resize_buffer(size_t new_buffer_size) {
115+
metrics_q_.resize(new_buffer_size);
116+
}
117+
118+
void GrpcMetricsExporter::on_metrics_export_complete(::grpc::Status status) {
119+
in_flight_ = false;
120+
if (status.ok()) {
121+
metrics_q_.pop();
122+
export_current();
123+
}
124+
}
125+
126+
} // namespace grpc
127+
} // namespace nsolid
128+
} // namespace node

0 commit comments

Comments
 (0)