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
111 changes: 99 additions & 12 deletions v2/lib/telemetry.zig
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,10 @@ pub const Region = extern struct {
pub const Info = extern struct {
/// The port to listen on for the prometheus client.
port: u16,
/// The most verbose level enabled by any filter in `log_filters_encoded`.
/// Writers gate on this; the exact (service, scope) filter is still applied
/// by the telemetry service in `log.streamLogs`.
max_log_level: log.Level,
/// The length of the encoded log filters byte string.
log_filters_len: u32,

Expand All @@ -65,6 +69,19 @@ pub const Region = extern struct {
/// A histogram with N buckets consumes `3 * N + 6` elements.
histogram_data_len: u32,

/// The padded size of the region header, with no trailing allocations.
/// Only the `*_len` and `service_count` fields affect layout. `port` and
/// `max_log_level` are configuration, and arbitrary here.
const header_padded_size = Info.regionSize(.{
.port = 0,
.max_log_level = .trace,
.log_filters_len = 0,
.service_count = 0,
.id_mem_len = 0,
.gauges_len = 0,
.histogram_data_len = 0,
});

/// NOTE: keep in sync with `Region.getSlices`.
pub fn regionSize(self: Info) usize {
var size: usize = 0;
Expand Down Expand Up @@ -112,6 +129,7 @@ pub const Region = extern struct {
pub fn info(self: InitParams) Info {
return .{
.port = self.port,
.max_log_level = log.maxLevelEncoded(self.log_filters_encoded),
.log_filters_len = @intCast(self.log_filters_encoded.len),

.service_count = self.service_count,
Expand Down Expand Up @@ -153,15 +171,7 @@ pub const Region = extern struct {
const buf: []align(@alignOf(u64)) u8 = trailing: {
const ptr: [*]align(@alignOf(u64)) u8 = @ptrCast(self);
const full: []align(@alignOf(u64)) u8 = ptr[0..self.info.regionSize()];
const header_padded_size = comptime Info.regionSize(.{
.port = 0,
.log_filters_len = 0,
.service_count = 0,
.id_mem_len = 0,
.gauges_len = 0,
.histogram_data_len = 0,
});
break :trailing full[header_padded_size..];
break :trailing full[Info.header_padded_size..];
};
var seek: usize = 0;
const log_filters =
Expand Down Expand Up @@ -218,7 +228,10 @@ pub const Region = extern struct {

std.debug.assert(name.len <= log.MessageStream.Name.MAX_LEN); // see `stream.name.init`
stream.name.init(name);
return .{ .sink = .{ .swap_buffer = &stream.swap_buffer } };
return .{
.sink = .{ .swap_buffer = &stream.swap_buffer },
.max_level = self.info.max_log_level,
};
}

/// Low-level helper for registering metrics.
Expand All @@ -240,11 +253,15 @@ pub const Region = extern struct {
pub fn Logger(comptime scope_str: []const u8) type {
return struct {
sink: log.MessageSink,
/// See `Region.Info.max_log_level`. Scope-independent, so it survives
/// `withScope`/`from` unchanged.
max_level: log.Level,

const LoggerSelf = @This();

pub const scope = scope_str;

pub const noop: LoggerSelf = .{ .sink = .noop };
pub const noop: LoggerSelf = .{ .sink = .noop, .max_level = .trace };

pub fn from(logger: anytype) LoggerSelf {
const LoggerOther = Logger(@TypeOf(logger).scope);
Expand All @@ -255,7 +272,7 @@ pub fn Logger(comptime scope_str: []const u8) type {
self: LoggerSelf,
comptime new_scope: []const u8,
) Logger(new_scope) {
return .{ .sink = self.sink };
return .{ .sink = self.sink, .max_level = self.max_level };
}

pub fn fatal(self: LoggerSelf) Entry(0) {
Expand Down Expand Up @@ -393,6 +410,13 @@ pub fn Logger(comptime scope_str: []const u8) type {
comptime fmt_str: []const u8,
args: anytype,
) void {
// `max_level` bounds what any filter in this process can enable, so
// `streamLogs` would drop this message regardless of service/scope.
// Bailing here avoids rendering it three times (twice to count length
// in `computeHeader`, once in `Message.write`) and avoids spending
// swap buffer on bytes nobody reads.
if (self.level.order(self.logger.max_level) == .gt) return;

switch (self.level) {
inline else => |ilevel| {
tracy.print(@tagName(ilevel) ++ ": " ++ fmt_str, args);
Expand Down Expand Up @@ -460,6 +484,69 @@ pub fn Logger(comptime scope_str: []const u8) type {
};
}

test "logf drops entries more verbose than max_level" {
var buf: [4096]u8 = undefined;

for (std.enums.values(log.Level)) |max_level| {
for (std.enums.values(log.Level)) |level| {
var fbw: std.Io.Writer = .fixed(&buf);
const logger: Logger("gate") = .{
.sink = .{ .writer = &fbw },
.max_level = max_level,
};
logger.entry(level).log("message");

if (level.order(max_level) == .gt) {
// Nothing may reach the sink; `streamLogs` would have dropped it anyway.
try std.testing.expectEqual(0, fbw.buffered().len);
continue;
}

// Everything at or below `max_level` is written whole and unaltered.
var fbr: std.Io.Reader = .fixed(fbw.buffered());
const header = try fbr.takeStruct(log.Message.Header, endian);
const slices = header.getSlicesFromFixedBuffer(&fbr) orelse
return error.TestExpectedNonNull;
try std.testing.expectEqual(.valid, header.magic);
try std.testing.expectEqual(level, header.level);
try std.testing.expectEqualStrings("gate", slices.scope);
try std.testing.expectEqualStrings("message", slices.msg);
try std.testing.expectEqual(0, fbr.bufferedLen());
}
}
}

test "max_level survives withScope and from" {
const base: Logger("base") = .{ .sink = .noop, .max_level = .debug };
try std.testing.expectEqual(.debug, base.withScope("scoped").max_level);
try std.testing.expectEqual(.debug, Logger("other").from(base).max_level);
}

test "acquireLogger takes max_level from the region's filters" {
const gpa = std.testing.allocator;

// `.warn` rather than the most verbose level, so that a logger which never consulted
// the region's filters would not pass.
const params: Region.InitParams = .{
.port = 0,
.log_filters_encoded = log.Filter.parseListStrLitIntoBinary(.warn, "").?,
.service_count = 1,
.id_mem_len = 0,
.gauges_len = 0,
.histogram_data_len = 0,
};
try std.testing.expectEqual(.warn, params.info().max_log_level);

// NOTE: dominated by the single `log.MessageStream`'s swap buffer; only the name
// and the encoded filters are ever touched here.
const region_buf = try gpa.alignedAlloc(u8, .of(u64), params.info().regionSize());
defer gpa.free(region_buf);

const region: *Region = @ptrCast(region_buf.ptr);
region.init(params);
try std.testing.expectEqual(.warn, region.acquireLogger("service", "scope").max_level);
}

pub const Counter = struct {
value: *std.atomic.Value(u64),

Expand Down
109 changes: 109 additions & 0 deletions v2/lib/telemetry/log.zig
Original file line number Diff line number Diff line change
Expand Up @@ -565,6 +565,29 @@ pub const Filter = struct {
}
};

/// Iterates over a byte string written by `parseListAndWriteBinary`.
pub const Iterator = struct {
fbr: std.Io.Reader,

pub fn init(encoded: []const u8) Iterator {
return .{ .fbr = .fixed(encoded) };
}

pub const NextError = error{InvalidFilter};

/// Returns `null` once `encoded` has been fully consumed. The `service` & `scope`
/// of the returned filter point into `encoded`.
///
/// Returns `error.InvalidFilter` for a truncated header, or a header whose
/// service/scope bytes are missing; the iterator should not be used afterwards.
pub fn next(self: *Iterator) NextError!?Filter {
if (self.fbr.bufferedLen() == 0) return null;
const header = self.fbr.takeStruct(Header, tel.endian) catch
return error.InvalidFilter;
return header.getFilterFromFixedReader(&self.fbr) orelse return error.InvalidFilter;
}
};

pub fn format(self: Filter, w: *std.Io.Writer) std.Io.Writer.Error!void {
if (self.service) |service| try w.writeAll(service);
if (self.scope) |scope| {
Expand Down Expand Up @@ -799,6 +822,92 @@ test EntryValueFmt {
});
}

/// The most verbose level that any filter in `encoded` can enable, across every
/// service and scope. `streamLogs` is guaranteed to drop anything more verbose
/// than this, so writers may skip encoding such messages entirely.
///
/// Returns `.trace` (gate fully open) for empty or malformed input, leaving the
/// diagnostic to the telemetry service, which rejects both explicitly.
pub fn maxLevelEncoded(encoded: []const u8) Level {
// NOTE: load-bearing; without it an empty list would fall through to `.fatal`,
// which is the strictest gate rather than the most open one.
if (encoded.len == 0) return .trace;
var max: Level = .fatal;
var filters: Filter.Iterator = .init(encoded);
while (filters.next() catch return .trace) |filter| {
if (filter.level.order(max) == .gt) max = filter.level;
}
return max;
}

test maxLevelEncoded {
@setEvalBranchQuota(16000); // encoding the filter lists below runs at comptime

// No filters at all leaves the gate open; the telemetry service is what rejects this.
try std.testing.expectEqual(.trace, maxLevelEncoded(""));

// The default filter counts towards the maximum.
try std.testing.expectEqual(.debug, maxLevelEncoded(
comptime Filter.parseListStrLitIntoBinary(.debug, "replay=error").?,
));

// So does any other filter, whichever position it holds in the list.
try std.testing.expectEqual(.trace, maxLevelEncoded(
comptime Filter.parseListStrLitIntoBinary(.fatal, "replay:main=trace,gossip=error").?,
));
try std.testing.expectEqual(.trace, maxLevelEncoded(
comptime Filter.parseListStrLitIntoBinary(.fatal, "gossip=error,replay:main=trace").?,
));

// Truncated input keeps the gate open rather than silently dropping messages, for both
// a partial header and a header whose service/scope bytes are missing.
{
const encoded = comptime Filter.parseListStrLitIntoBinary(.err, "replay:main=warn").?;
try std.testing.expectEqual(.warn, maxLevelEncoded(encoded));
try std.testing.expectEqual(.trace, maxLevelEncoded(encoded[0 .. encoded.len - 1]));
try std.testing.expectEqual(
.trace,
maxLevelEncoded(encoded[0 .. @sizeOf(Filter.Header) + 1]),
);
}

// The invariant all of the above serves: the result must bound every level `streamLogs`
// can select, otherwise the writer-side gate drops messages that it was going to emit,
// and they disappear with no diagnostic.
{
const encoded = comptime Filter.parseListStrLitIntoBinary(
.err,
"replay:main=trace,replay=debug,gossip:pull=info,accountsdb=warn",
).?;
const max = maxLevelEncoded(encoded);

// Decode into the sorted list `streamLogs` is given; see `services/telemetry.zig`.
var filters_buffer: [8]Filter = undefined;
var filters: std.ArrayList(Filter) = .initBuffer(&filters_buffer);
var iter: Filter.Iterator = .init(encoded);
while (try iter.next()) |filter| try filters.appendBounded(filter);
std.sort.block(Filter, filters.items, {}, Filter.sortLessThanInverted);

// Every service & scope named by the list, plus pairs that fall through to a
// broader filter or to the default.
for ([_][2][]const u8{
.{ "replay", "main" },
.{ "replay", "other" },
.{ "gossip", "pull" },
.{ "gossip", "push" },
.{ "accountsdb", "manager" },
.{ "unlisted", "scope" },
}) |pair| {
const index = Filter.findClosestFilter(.{
.filters = filters.items,
.service = pair[0],
.scope = pair[1],
});
try std.testing.expect(filters.items[index].level.order(max) != .gt);
}
}
}

pub fn streamLogs(
params: struct {
output: *std.Io.Writer,
Expand Down
2 changes: 1 addition & 1 deletion v2/lib/telemetry/tests/TestLogStore.zig
Original file line number Diff line number Diff line change
Expand Up @@ -244,7 +244,7 @@ pub fn deinit(self: *TestLogStore) void {

/// Returns a logger that writes records into this store.
pub fn logger(self: *TestLogStore, comptime scope: []const u8) tel.Logger(scope) {
return .{ .sink = .{ .writer = &self.state.writer } };
return .{ .sink = .{ .writer = &self.state.writer }, .max_level = .trace };
}

/// Returned records and iterators are invalidated by the next log, reset, or deinit.
Expand Down
10 changes: 2 additions & 8 deletions v2/services/telemetry.zig
Original file line number Diff line number Diff line change
Expand Up @@ -57,14 +57,8 @@ pub fn serviceMain(runner: lib.runner.Connection, ro: ReadOnly, rw: ReadWrite) !
var filters_buffer: [4096]api.log.Filter = undefined;
const filters: []const api.log.Filter = filters: {
var filters: std.ArrayList(api.log.Filter) = .initBuffer(&filters_buffer);
var fbr: std.Io.Reader = .fixed(region.getSlices().log_filters_encoded);
while (fbr.bufferedLen() != 0) {
const header = try fbr.takeStruct(api.log.Filter.Header, api.endian);
const filter = header.getFilterFromFixedReader(&fbr) orelse {
return error.InvalidFilterHeader;
};
try filters.appendBounded(filter);
}
var iter: api.log.Filter.Iterator = .init(region.getSlices().log_filters_encoded);
while (try iter.next()) |filter| try filters.appendBounded(filter);
std.sort.block(api.log.Filter, filters.items, {}, api.log.Filter.sortLessThanInverted);
if (filters.items.len == 0 or !filters.getLast().isLevelOnly()) {
std.log.err(
Expand Down
Loading