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
20 changes: 13 additions & 7 deletions spectator/publisher.cc
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
#include "publisher.h"
#include "logger.h"
#include <array>
#include <fmt/format.h>

namespace spectator {
Expand Down Expand Up @@ -37,7 +38,9 @@ SpectatordPublisher::SpectatordPublisher(absl::string_view endpoint,
}

void SpectatordPublisher::setup_nop_sender() {
sender_ = [this](std::string_view msg) { logger_->trace("{}", msg); };
sender_ = [this](std::string_view prefix, std::string_view value) {
logger_->trace("{}{}", prefix, value);
};
}

void SpectatordPublisher::local_reconnect(absl::string_view path) {
Expand Down Expand Up @@ -76,8 +79,9 @@ void SpectatordPublisher::setup_unix_domain(absl::string_view path) {
buffer_.clear();
};

sender_ = [this](std::string_view msg) {
buffer_.append(msg);
sender_ = [this](std::string_view prefix, std::string_view value) {
buffer_.append(prefix);
buffer_.append(value);
const auto now = std::chrono::steady_clock::now();
const bool should_flush = buffer_.length() >= bytes_to_buffer_ ||
now - last_flush_time_ >= flush_interval_;
Expand Down Expand Up @@ -126,14 +130,16 @@ void SpectatordPublisher::udp_reconnect(
void SpectatordPublisher::setup_udp(absl::string_view host_port) {
auto endpoint = resolve_host_port(io_context_, host_port);
udp_reconnect(endpoint);
sender_ = [endpoint, this](std::string_view msg) {
sender_ = [endpoint, this](std::string_view prefix, std::string_view value) {
std::array<asio::const_buffer, 2> buffers{asio::buffer(prefix.data(), prefix.size()),
asio::buffer(value.data(), value.size())};
for (auto i = 0; i < 3; ++i) {
try {
udp_socket_.send(asio::buffer(msg));
logger_->trace("Sent (udp): {}", msg);
udp_socket_.send(buffers);
logger_->trace("Sent (udp): {}{}", prefix, value);
break;
} catch (std::exception& e) {
logger_->warn("Unable to send {} - attempt {}/3", msg, i);
logger_->warn("Unable to send {}{} - attempt {}/3", prefix, value, i);
udp_reconnect(endpoint);
}
}
Expand Down
5 changes: 3 additions & 2 deletions spectator/publisher.h
Original file line number Diff line number Diff line change
Expand Up @@ -24,11 +24,12 @@ class SpectatordPublisher {
std::shared_ptr<spdlog::logger> logger = DefaultLogger());
SpectatordPublisher(const SpectatordPublisher&) = delete;

void send(std::string_view measurement) { sender_(measurement); };
// Sends a measured value for the given metric name prefix.
void send(std::string_view prefix, std::string_view value) { sender_(prefix, value); };
Comment thread
jasonk000 marked this conversation as resolved.
void flush() { flusher_(); };

protected:
using sender_fun = std::function<void(std::string_view)>;
using sender_fun = std::function<void(std::string_view, std::string_view)>;
sender_fun sender_;
using flusher_fun = std::function<void()>;
flusher_fun flusher_ = []() {};
Expand Down
49 changes: 18 additions & 31 deletions spectator/stateless_meters.h
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
#include <cassert>
#include <charconv>
#include <cmath>
#include <memory>
#include "id.h"
#include "absl/strings/str_cat.h"
#include "absl/time/time.h"
Expand Down Expand Up @@ -50,14 +51,6 @@ inline std::string create_prefix(const Id& id, std::string_view type_name) {
return res;
}

// Single thread-local send buffer shared across all StatelessMeter instantiations.
// Non-template so all Pub types resolve to the same storage slot per thread.
// Not re-entrant: callers must complete send() before the buffer is safe to reuse.
inline std::string& tl_send_buf() {
thread_local std::string buf;
return buf;
}

template <typename T>
T restrict(T amount, T min, T max) {
auto r = amount;
Expand Down Expand Up @@ -88,19 +81,15 @@ class StatelessMeter {
protected:
void send(double value) {
ensure_prefix();
auto& tl_msg = detail::tl_send_buf();
tl_msg.assign(value_prefix_);

// Early exit: match absl::StrFormat("%f") behaviour for special values.
// These are passed straight through as literals -- no buffer needed.
if (std::isnan(value)) {
tl_msg.append("nan");
publisher_->send(tl_msg);
publisher_->send(value_prefix_, std::string_view("nan"));
return;
}

if (std::isinf(value)) {
tl_msg.append(value > 0 ? "inf" : "-inf");
publisher_->send(tl_msg);
publisher_->send(value_prefix_, std::string_view(value > 0 ? "inf" : "-inf"));
return;
}

Expand All @@ -111,30 +100,28 @@ class StatelessMeter {
auto [ptr, ec] = std::to_chars(num_buf, num_buf + sizeof(num_buf), value,
std::chars_format::fixed);
if (ec == std::errc{}) {
tl_msg.append(num_buf, ptr);
} else {
// Fallback for subnormal values, which require up to 1076 chars in fixed
// notation. NaN/Inf are handled above, so this branch is subnormals only.
auto off = tl_msg.size();
tl_msg.resize(off + detail::kMaxFixedDoubleLen);
auto [heap_ptr, heap_ec] = std::to_chars(tl_msg.data() + off,
tl_msg.data() + tl_msg.size(), value,
std::chars_format::fixed);
assert(heap_ec == std::errc{});
tl_msg.resize(static_cast<size_t>(heap_ptr - tl_msg.data()));
publisher_->send(value_prefix_, std::string_view(num_buf, static_cast<size_t>(ptr - num_buf)));
return;
}
publisher_->send(tl_msg);

// Slow path: didn't fit in 64 bytes. Only subnormal values hit this,
// which require up to 1076 chars in fixed notation (NaN/Inf are
// handled above). Heap-allocate once per thread to avoid inflating
// stack frames on the hot path.
thread_local auto big = std::make_unique<char[]>(detail::kMaxFixedDoubleLen);
auto [heap_ptr, heap_ec] = std::to_chars(big.get(), big.get() + detail::kMaxFixedDoubleLen,
value, std::chars_format::fixed);
assert(heap_ec == std::errc{});
publisher_->send(value_prefix_,
std::string_view(big.get(), static_cast<size_t>(heap_ptr - big.get())));
}

void send_uint(uint64_t value) {
ensure_prefix();
char num_buf[24];
auto [ptr, ec] = std::to_chars(num_buf, num_buf + sizeof(num_buf), value);
assert(ec == std::errc{});
auto& tl_msg = detail::tl_send_buf();
tl_msg.assign(value_prefix_);
tl_msg.append(num_buf, ptr);
publisher_->send(tl_msg);
publisher_->send(value_prefix_, std::string_view(num_buf, static_cast<size_t>(ptr - num_buf)));
}

private:
Expand Down
4 changes: 3 additions & 1 deletion spectator/test_publisher.h
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,9 @@
namespace spectator {
class TestPublisher {
public:
void send(std::string_view msg) { messages.emplace_back(msg); }
void send(std::string_view prefix, std::string_view value) {
messages.emplace_back(std::string(prefix) + std::string(value));
}
std::vector<std::string> SentMessages() { return messages; }
void Reset() { messages.clear(); }

Expand Down
Loading