Skip to content
Merged
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
68 changes: 28 additions & 40 deletions src/core/statsmanager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -301,9 +301,8 @@ class StatsManager::Private {
int subscriptionLinger;
int reportInterval;
std::unique_ptr<ZmqSocket> sock;
QString prometheusPrefix;
ffi::CommonMetrics *commonMetrics;
ffi::PrometheusServer *prometheusServer;
std::shared_ptr<CommonMetrics> commonMetrics;
size_t commonMetricsRegistrationId;
QHash<QByteArray, uint32_t> routeActivity;
QHash<QByteArray, ConnectionInfo *> connectionInfoById;
QHash<QByteArray, QSet<ConnectionInfo *>> connectionInfoByRoute;
Expand Down Expand Up @@ -345,8 +344,7 @@ class StatsManager::Private {
subscriptionTtl(60 * 1000),
subscriptionLinger(60 * 1000),
reportInterval(10 * 1000),
commonMetrics(nullptr),
prometheusServer(nullptr),
commonMetricsRegistrationId(0),
currentConnectionInfoRefreshBucket(0),
currentSubscriptionRefreshBucket(0),
wheel(TimerWheel((_connectionsMax * 2) + _subscriptionsMax)) {
Expand Down Expand Up @@ -374,8 +372,8 @@ class StatsManager::Private {
}

~Private() {
ffi::prometheus_server_destroy(prometheusServer);
ffi::statsmanager_commonmetrics_destroy(commonMetrics);
if (commonMetrics)
commonMetrics->unregisterInstance(commonMetricsRegistrationId);

qDeleteAll(connectionInfoById);

Expand Down Expand Up @@ -409,35 +407,17 @@ class StatsManager::Private {
return true;
}

bool setPrometheusPort(const QString &portStr) {
assert(!commonMetrics && !prometheusServer);

commonMetrics = ffi::statsmanager_commonmetrics_create(prometheusPrefix.toUtf8().data());
if (!commonMetrics)
return false;

const ffi::PrometheusRegistry *registry =
ffi::statsmanager_commonmetrics_registry(commonMetrics);

const char *error = nullptr;
prometheusServer = ffi::prometheus_server_create(portStr.toUtf8().data(), registry, &error);
if (!prometheusServer) {
log_error("prometheus_server_create: %s", error);
ffi::prometheus_server_error_destroy(error);
ffi::statsmanager_commonmetrics_destroy(commonMetrics);
commonMetrics = nullptr;
return false;
}

return true;
void setCommonMetrics(std::shared_ptr<CommonMetrics> cm) {
assert(!commonMetrics);
commonMetrics = std::move(cm);
commonMetricsRegistrationId = commonMetrics->registerInstance();
}

void combinedReportChanged() {
if (commonMetrics) {
ffi::statsmanager_commonmetrics_update(
commonMetrics, combinedReport.requestsReceived, combinedReport.connectionsMax,
combinedReport.connectionsMinutes, combinedReport.messagesReceived,
combinedReport.messagesSent);
commonMetrics->update(commonMetricsRegistrationId, combinedReport.requestsReceived,
combinedReport.connectionsMax, combinedReport.connectionsMinutes,
combinedReport.messagesReceived, combinedReport.messagesSent);
}
}

Expand Down Expand Up @@ -1332,17 +1312,25 @@ StatsManager::CommonMetrics::create(const QString &prefix) {
ffi::CommonMetrics *handle = ffi::statsmanager_commonmetrics_create(prefix.toUtf8().data());
if (!handle)
return nullptr;
return std::unique_ptr<StatsManager::CommonMetrics>(new StatsManager::CommonMetrics(handle));
return std::unique_ptr<CommonMetrics>(new CommonMetrics(handle));
}

const ffi::PrometheusRegistry *StatsManager::CommonMetrics::registry() const {
return ffi::statsmanager_commonmetrics_registry(inner_);
}

void StatsManager::CommonMetrics::update(uint32_t requestReceived, uint32_t connectionConnected,
uint32_t connectionMinute, uint32_t messageReceived,
uint32_t messageSent) {
ffi::statsmanager_commonmetrics_update(inner_, requestReceived, connectionConnected,
size_t StatsManager::CommonMetrics::registerInstance() {
return ffi::statsmanager_commonmetrics_register(inner_);
}

void StatsManager::CommonMetrics::unregisterInstance(size_t id) {
ffi::statsmanager_commonmetrics_unregister(inner_, id);
}

void StatsManager::CommonMetrics::update(size_t id, uint32_t requestReceived,
uint32_t connectionConnected, uint32_t connectionMinute,
uint32_t messageReceived, uint32_t messageSent) {
ffi::statsmanager_commonmetrics_update(inner_, id, requestReceived, connectionConnected,
connectionMinute, messageReceived, messageSent);
}

Expand Down Expand Up @@ -1388,9 +1376,9 @@ void StatsManager::setReportInterval(int secs) {

void StatsManager::setOutputFormat(Format format) { d->outputFormat = format; }

bool StatsManager::setPrometheusPort(const QString &port) { return d->setPrometheusPort(port); }

void StatsManager::setPrometheusPrefix(const QString &prefix) { d->prometheusPrefix = prefix; }
void StatsManager::setCommonMetrics(std::shared_ptr<CommonMetrics> commonMetrics) {
d->setCommonMetrics(std::move(commonMetrics));
}

void StatsManager::addActivity(const QByteArray &routeId, uint32_t count) {
if (d->routeActivity.contains(routeId))
Expand Down
13 changes: 9 additions & 4 deletions src/core/statsmanager.h
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,8 @@
#include "rust/bindings.h"
#include "stats.h"
#include <boost/signals2.hpp>
#include <cstddef>
#include <memory>

class QHostAddress;

Expand All @@ -42,7 +44,9 @@ class StatsManager {
enum Format { TnetStringFormat, JsonFormat };

/// RAII wrapper around the Rust-backed CommonMetrics object, which holds the prometheus
/// registry and all metric handles.
/// registry and all metric handles. Multiple StatsManager instances may share one
/// CommonMetrics via std::shared_ptr; each registers itself to obtain a per-instance ID
/// that is passed to update().
class CommonMetrics {
public:
~CommonMetrics();
Expand All @@ -53,7 +57,9 @@ class StatsManager {
static std::unique_ptr<CommonMetrics> create(const QString &prefix);

const ffi::PrometheusRegistry *registry() const;
void update(uint32_t requestReceived, uint32_t connectionConnected,
size_t registerInstance();
void unregisterInstance(size_t id);
void update(size_t id, uint32_t requestReceived, uint32_t connectionConnected,
uint32_t connectionMinute, uint32_t messageReceived, uint32_t messageSent);

private:
Expand All @@ -77,8 +83,7 @@ class StatsManager {
void setSubscriptionLinger(int secs);
void setReportInterval(int secs);
void setOutputFormat(Format format);
bool setPrometheusPort(const QString &port);
void setPrometheusPrefix(const QString &prefix);
void setCommonMetrics(std::shared_ptr<CommonMetrics> commonMetrics);

// RouteId may be empty for non-identified route

Expand Down
126 changes: 99 additions & 27 deletions src/core/statsmanager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,16 @@

use crate::core::prometheus::try_register_process_collector;
use prometheus::{IntCounter, IntGauge};
use slab::Slab;

#[derive(Default)]
struct PrevValues {
request_received: u32,
connection_connected: u32,
connection_minute: u32,
message_received: u32,
message_sent: u32,
}

/// Metrics needed by the `StatsManager` C++ class.
pub struct CommonMetrics {
Expand All @@ -25,10 +35,7 @@ pub struct CommonMetrics {
connection_minute: IntCounter,
message_received: IntCounter,
message_sent: IntCounter,
prev_request_received: u32,
prev_connection_minute: u32,
prev_message_received: u32,
prev_message_sent: u32,
instances: Slab<PrevValues>,
}

impl CommonMetrics {
Expand Down Expand Up @@ -88,45 +95,75 @@ impl CommonMetrics {
connection_minute,
message_received,
message_sent,
prev_request_received: 0,
prev_connection_minute: 0,
prev_message_received: 0,
prev_message_sent: 0,
instances: Slab::new(),
}
}

fn register(&mut self) -> usize {
self.instances.insert(PrevValues::default())
}

fn unregister(&mut self, id: usize) {
let prev = self.instances.remove(id);
if prev.connection_connected > 0 {
self.connection_connected
.sub(prev.connection_connected as i64);
}
}

fn update(
&mut self,
id: usize,
request_received: u32,
connection_connected: u32,
connection_minute: u32,
message_received: u32,
message_sent: u32,
) {
let delta = request_received.saturating_sub(self.prev_request_received);
if delta > 0 {
self.request_received.inc_by(delta as u64);
self.prev_request_received = request_received;
// Compute deltas and update prev values in a scoped borrow so we can
// subsequently call methods on the rest of `self`.
let (req_delta, conn_delta, conn_min_delta, msg_recv_delta, msg_sent_delta) = {
let prev = &mut self.instances[id];

let req_delta = request_received.saturating_sub(prev.request_received);
let conn_delta = (connection_connected as i64) - (prev.connection_connected as i64);
let conn_min_delta = connection_minute.saturating_sub(prev.connection_minute);
let msg_recv_delta = message_received.saturating_sub(prev.message_received);
let msg_sent_delta = message_sent.saturating_sub(prev.message_sent);

prev.request_received = request_received;
prev.connection_connected = connection_connected;
prev.connection_minute = connection_minute;
prev.message_received = message_received;
prev.message_sent = message_sent;

(
req_delta,
conn_delta,
conn_min_delta,
msg_recv_delta,
msg_sent_delta,
)
};

if req_delta > 0 {
self.request_received.inc_by(req_delta as u64);
}

self.connection_connected.set(connection_connected as i64);
if conn_delta != 0 {
self.connection_connected.add(conn_delta);
}

let delta = connection_minute.saturating_sub(self.prev_connection_minute);
if delta > 0 {
self.connection_minute.inc_by(delta as u64);
self.prev_connection_minute = connection_minute;
if conn_min_delta > 0 {
self.connection_minute.inc_by(conn_min_delta as u64);
}

let delta = message_received.saturating_sub(self.prev_message_received);
if delta > 0 {
self.message_received.inc_by(delta as u64);
self.prev_message_received = message_received;
if msg_recv_delta > 0 {
self.message_received.inc_by(msg_recv_delta as u64);
}

let delta = message_sent.saturating_sub(self.prev_message_sent);
if delta > 0 {
self.message_sent.inc_by(delta as u64);
self.prev_message_sent = message_sent;
if msg_sent_delta > 0 {
self.message_sent.inc_by(msg_sent_delta as u64);
}
}
}
Expand Down Expand Up @@ -188,15 +225,49 @@ mod ffi {
&m.registry as *const prometheus::Registry as *const PrometheusRegistry
}

/// Update all metrics to the current totals. For counters the delta since the last call is
/// computed internally; `connection_connected` is a gauge and is set directly.
/// Register a new StatsManager instance and return an opaque ID for it. Pass this ID to
/// `statsmanager_commonmetrics_update` and `statsmanager_commonmetrics_unregister`.
///
/// # Safety
///
/// `m` must be a valid non-null pointer returned by `statsmanager_commonmetrics_create`.
#[no_mangle]
pub unsafe extern "C" fn statsmanager_commonmetrics_register(m: *mut CommonMetrics) -> usize {
let m = unsafe { m.as_mut().unwrap() };
m.register()
}

/// Unregister a StatsManager instance previously registered with
/// `statsmanager_commonmetrics_register`. After this call, `id` must not be passed to
/// `statsmanager_commonmetrics_update`.
///
/// # Safety
///
/// `m` must be a valid non-null pointer returned by `statsmanager_commonmetrics_create`.
/// `id` must be a value previously returned by `statsmanager_commonmetrics_register` on
/// the same instance that has not yet been unregistered.
#[no_mangle]
pub unsafe extern "C" fn statsmanager_commonmetrics_unregister(
m: *mut CommonMetrics,
id: usize,
) {
let m = unsafe { m.as_mut().unwrap() };
m.unregister(id);
}

/// Update all metrics to the current totals for the given StatsManager instance. For counters
/// the delta since the last call (per instance) is computed internally;
/// `connection_connected` is a gauge and is set directly.
///
/// # Safety
///
/// `m` must be a valid non-null pointer returned by `statsmanager_commonmetrics_create`.
/// `id` must be a value previously returned by `statsmanager_commonmetrics_register` on
/// the same instance that has not yet been unregistered.
#[no_mangle]
pub unsafe extern "C" fn statsmanager_commonmetrics_update(
m: *mut CommonMetrics,
id: usize,
request_received: u32,
connection_connected: u32,
connection_minute: u32,
Expand All @@ -206,6 +277,7 @@ mod ffi {
let m = unsafe { m.as_mut().unwrap() };

m.update(
id,
request_received,
connection_connected,
connection_minute,
Expand Down
15 changes: 11 additions & 4 deletions src/handler/handlerengine.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
#include "packet/retryrequestpacket.h"
#include "packet/statspacket.h"
#include "packet/wscontrolpacket.h"
#include "prometheus.h"
#include "publishformat.h"
#include "publishitem.h"
#include "publishlastids.h"
Expand Down Expand Up @@ -1120,6 +1121,8 @@ class HandlerEngine::Private {
std::unique_ptr<ZmqSocket> proxyStatsSock;
std::unique_ptr<ZmqValve> proxyStatsValve;
std::unique_ptr<SimpleHttpServer> controlHttpServer;
std::shared_ptr<StatsManager::CommonMetrics> commonMetrics;
std::unique_ptr<PrometheusServer> prometheusServer;
std::unique_ptr<StatsManager> stats;
std::unique_ptr<RateLimiter> publishLimiter;
std::unique_ptr<RateLimiter> updateLimiter;
Expand Down Expand Up @@ -1391,13 +1394,17 @@ class HandlerEngine::Private {
}

if (!config.prometheusPort.isEmpty()) {
stats->setPrometheusPrefix(config.prometheusPrefix);
commonMetrics = StatsManager::CommonMetrics::create(config.prometheusPrefix);

if (!stats->setPrometheusPort(config.prometheusPort)) {
log_error("unable to bind to prometheus port: %s",
qPrintable(config.prometheusPort));
QString promError;
prometheusServer = PrometheusServer::create(config.prometheusPort,
commonMetrics->registry(), &promError);
if (!prometheusServer) {
log_error("unable to bind to prometheus port: %s", qPrintable(promError));
return false;
}

stats->setCommonMetrics(commonMetrics);
}

if (!config.proxyStatsSpecs.isEmpty()) {
Expand Down
Loading
Loading