Skip to content

Commit 54f8e42

Browse files
committed
feat: validator monitor
Implements the missing validator monitor metrics from lodestar-ts. Note that there are some metrics intentionally omitted because they probably require even bigger changes on the zig side, this is just the minimal change required to track validator perf. The following metrics are excluded: - `validator_monitor_prev_epoch_on_chain_attester_hit_total` - `validator_monitor_prev_epoch_on_chain_attester_miss_total` - `validator_monitor_prev_epoch_on_chain_attester_correct_head_total` - `validator_monitor_prev_epoch_on_chain_attester_incorrect_head_total` - `validator_monitor_prev_epoch_on_chain_inclusion_distance` These will be updated on the TS side to run on `onceEveryEndOfEpoch` rather than inside `registerValidatorStatuses`, which is now zig native
1 parent 2145faa commit 54f8e42

13 files changed

Lines changed: 366 additions & 7 deletions

File tree

bench/state_transition/process_block.zig

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -529,6 +529,7 @@ fn runBenchmark(
529529
cached_state,
530530
block_slot,
531531
.{},
532+
null,
532533
);
533534
try cached_state.state.commit();
534535
try state_transition.buildSlashingsCacheFromStateIfNeeded(allocator, cached_state.state, &cached_state.slashings_cache);

bindings/napi/BeaconStateView.zig

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ const js_types = @import("./js_types.zig");
1717
const sszValueToNapiValue = @import("./to_napi_value.zig").sszValueToNapiValue;
1818
const numberSliceToNapiValue = @import("./to_napi_value.zig").numberSliceToNapiValue;
1919
const napi_io = @import("./io.zig");
20+
const validator_monitor = @import("./validator_monitor.zig");
2021

2122
/// Allocator used for all BeaconStateView instances.
2223
var gpa: std.heap.DebugAllocator(.{}) = .init;
@@ -1289,7 +1290,14 @@ pub fn processSlots(self: *const BeaconStateView, slot_arg: js.Number, options:
12891290
allocator.destroy(post_state);
12901291
}
12911292

1292-
try st.processSlots(allocator, napi_io.get(), post_state, slot_value, .{});
1293+
try st.processSlots(
1294+
allocator,
1295+
napi_io.get(),
1296+
post_state,
1297+
slot_value,
1298+
.{},
1299+
validator_monitor.get(),
1300+
);
12931301
return .{
12941302
.cached_state = post_state,
12951303
.pool_rc = pool.state.poolRc().ref(),
@@ -1318,7 +1326,14 @@ pub fn stateTransition(self: *const BeaconStateView, signed_block_bytes: js.Uint
13181326
const signed_block = try AnySignedBeaconBlock.deserialize(allocator, .full, fork_seq, bytes);
13191327
defer signed_block.deinit(allocator);
13201328

1321-
const post_state = try st.stateTransition(allocator, napi_io.get(), cached_state, signed_block, opts);
1329+
const post_state = try st.stateTransition(
1330+
allocator,
1331+
napi_io.get(),
1332+
cached_state,
1333+
signed_block,
1334+
opts,
1335+
validator_monitor.get(),
1336+
);
13221337
return .{
13231338
.cached_state = post_state,
13241339
.pool_rc = pool.state.poolRc().ref(),

bindings/napi/metrics.zig

Lines changed: 23 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -12,13 +12,34 @@ else
1212

1313
var initialized: bool = false;
1414

15+
const validator_monitor = @import("./validator_monitor.zig");
16+
1517
/// JS: metrics.init() → void
1618
pub fn init() !void {
1719
if (initialized) return;
1820
try state_transition.metrics.init(allocator, napi_io.get(), .{});
1921
initialized = true;
2022
}
2123

24+
/// JS: metrics.registerLocalValidator(index) → void
25+
///
26+
/// Adds a validator index to the process-wide validator monitor so that
27+
/// metrics are recorded for it on every epoch transition.
28+
pub fn registerLocalValidator(index: js.Number) !void {
29+
const value = try index.toI64();
30+
if (value < 0) return error.InvalidValidatorIndex;
31+
try validator_monitor.get().registerLocalValidator(napi_io.get(), @intCast(value));
32+
}
33+
34+
/// JS: metrics.unregisterLocalValidator(index) → void
35+
///
36+
/// Prunes a validator index from the process-wide validator monitor.
37+
pub fn unregisterLocalValidator(index: js.Number) !void {
38+
const value = try index.toI64();
39+
if (value < 0) return error.InvalidValidatorIndex;
40+
validator_monitor.get().unregisterLocalValidator(napi_io.get(), @intCast(value));
41+
}
42+
2243
/// JS: metrics.scrapeMetrics() → string
2344
pub fn scrapeMetrics() !js.String {
2445
var aw: std.Io.Writer.Allocating = .init(allocator);
@@ -29,7 +50,8 @@ pub fn scrapeMetrics() !js.String {
2950
}
3051

3152
pub fn deinit() void {
53+
validator_monitor.deinit();
3254
if (!initialized) return;
33-
state_transition.metrics.state_transition.deinit();
55+
state_transition.metrics.deinit();
3456
initialized = false;
3557
}
Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,31 @@
1+
//! Process-wide validator monitor shared by the NAPI bindings.
2+
//!
3+
//! This file is intentionally NOT exported as a JS module in `root.zig`:
4+
//! it only holds native state. JS interacts with it through
5+
//! `metrics.registerLocalValidator()` and the metrics scraped via
6+
//! `metrics.scrapeMetrics()`.
7+
8+
const std = @import("std");
9+
const builtin = @import("builtin");
10+
const state_transition = @import("state_transition");
11+
12+
var gpa: std.heap.DebugAllocator(.{}) = .init;
13+
const allocator = if (builtin.mode == .Debug)
14+
gpa.allocator()
15+
else
16+
std.heap.c_allocator;
17+
18+
/// Fed by `BeaconStateView.processSlots`/`stateTransition` on every epoch
19+
/// transition. Only records metrics for validators registered via
20+
/// `metrics.registerLocalValidator()`.
21+
var monitor = state_transition.ValidatorMonitor.init(allocator);
22+
23+
/// Returns the process-wide validator monitor.
24+
pub fn get() *state_transition.ValidatorMonitor {
25+
return &monitor;
26+
}
27+
28+
/// Frees all monitor state. Meant to be called once on module cleanup.
29+
pub fn deinit() void {
30+
monitor.deinit();
31+
}

bindings/src/index.d.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -384,6 +384,8 @@ declare const bindings: {
384384
metrics: {
385385
init: () => void;
386386
scrapeMetrics: () => string;
387+
registerLocalValidator: (index: number) => void;
388+
unregisterLocalValidator: (index: number) => void;
387389
};
388390
BeaconStateView: typeof BeaconStateView;
389391
};

bindings/src/metrics.d.ts

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,3 +3,16 @@ export declare function init(): void;
33

44
/** Scrape native state-transition metrics in Prometheus text format. */
55
export declare function scrapeMetrics(): string;
6+
7+
/**
8+
* Register a validator index with the native validator monitor. Metrics
9+
* are recorded for registered validators on every epoch transition.
10+
*/
11+
export declare function registerLocalValidator(index: number): void;
12+
13+
/**
14+
* Remove a validator index from the native validator monitor, so its
15+
* `validator_monitor_*` status metrics stop being recorded. Mirrors the pruning
16+
* of stale registrations in lodestar's validator monitor.
17+
*/
18+
export declare function unregisterLocalValidator(index: number): void;

bindings/src/metrics.js

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,3 +4,5 @@ const native = bindings.metrics;
44

55
export const init = native.init;
66
export const scrapeMetrics = native.scrapeMetrics;
7+
export const registerLocalValidator = native.registerLocalValidator;
8+
export const unregisterLocalValidator = native.unregisterLocalValidator;
Lines changed: 194 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,194 @@
1+
//! Consumes the per-validator status data produced by the epoch transition
2+
//! and records validator metrics for the validators registered with `registerLocalValidator`.
3+
4+
const std = @import("std");
5+
6+
const types = @import("consensus_types");
7+
const metrics = @import("metrics.zig");
8+
const attester_status = @import("utils/attester_status.zig");
9+
const hasMarkers = attester_status.hasMarkers;
10+
11+
const Allocator = std.mem.Allocator;
12+
const Epoch = types.primitive.Epoch.Type;
13+
const ValidatorIndex = types.primitive.ValidatorIndex.Type;
14+
15+
pub const ValidatorMonitor = @This();
16+
17+
allocator: Allocator,
18+
/// Unordered list of validators that require additional monitoring.
19+
validators: std.AutoArrayHashMapUnmanaged(ValidatorIndex, void),
20+
/// Prevents registering statuses for the same epoch twice.
21+
/// processEpoch() may be run more than once for the same epoch.
22+
last_registered_status_epoch: ?Epoch,
23+
24+
pub fn init(allocator: Allocator) ValidatorMonitor {
25+
return .{
26+
.allocator = allocator,
27+
.validators = .empty,
28+
.last_registered_status_epoch = null,
29+
};
30+
}
31+
32+
pub fn deinit(self: *ValidatorMonitor) void {
33+
self.validators.deinit(self.allocator);
34+
self.* = undefined;
35+
}
36+
37+
/// Adds a validator to the list of monitored validators.
38+
///
39+
/// Registering an already-monitored validator is a no-op.
40+
pub fn registerLocalValidator(self: *ValidatorMonitor, io: std.Io, index: ValidatorIndex) !void {
41+
try self.validators.put(self.allocator, index, {});
42+
}
43+
44+
/// Prunes a validator from the list of monitored validators.
45+
///
46+
/// Unregistering an unknown validator is a no-op.
47+
pub fn unregisterLocalValidator(self: *ValidatorMonitor, io: std.Io, index: ValidatorIndex) void {
48+
_ = self.validators.swapRemove(index);
49+
}
50+
51+
/// Registers the per-validator statuses produced by one epoch transition and
52+
/// records metrics for all monitored validators.
53+
///
54+
/// `flags` are the packed attester flags of `EpochTransitionCache` (see `utils/attester_status.zig`).
55+
/// `balances` is optional; when present the total balance of all monitored
56+
/// validators is reported.
57+
///
58+
/// NOTE: Gossip-derived metrics (inclusion distance, attester hit/miss, correct head)
59+
/// are recorded by the lodestar-ts validator monitor, which observes attestations
60+
/// on the network; the state transition cannot see them post-altair.
61+
/// TODO(bing): port the rest of validator monitor, this is just a minimal
62+
/// impl for stf integration
63+
pub fn registerValidatorStatuses(
64+
self: *ValidatorMonitor,
65+
current_epoch: Epoch,
66+
flags: []const u8,
67+
balances: ?[]const u64,
68+
) void {
69+
70+
// Prevent registering statuses for the same epoch twice.
71+
if (self.last_registered_status_epoch) |last_registered_status_epoch|
72+
if (current_epoch <= last_registered_status_epoch) return;
73+
74+
self.last_registered_status_epoch = current_epoch;
75+
76+
// There won't be any validator activity in epoch -1.
77+
if (current_epoch == 0) return;
78+
79+
const vm = &metrics.validator_monitor;
80+
81+
// Track total balance instead of per-validator balance to reduce metric cardinality.
82+
var total_balance: u64 = 0;
83+
84+
for (self.validators.keys()) |index| {
85+
// The monitored validator may not be in the state yet.
86+
if (index >= flags.len) continue;
87+
88+
const flag = flags[index];
89+
90+
if (hasMarkers(flag, attester_status.FLAG_PREV_SOURCE_ATTESTER)) {
91+
vm.prev_epoch_on_chain_source_attester_hit.incr();
92+
} else {
93+
vm.prev_epoch_on_chain_source_attester_miss.incr();
94+
}
95+
if (hasMarkers(flag, attester_status.FLAG_PREV_HEAD_ATTESTER)) {
96+
vm.prev_epoch_on_chain_head_attester_hit.incr();
97+
} else {
98+
vm.prev_epoch_on_chain_head_attester_miss.incr();
99+
}
100+
if (hasMarkers(flag, attester_status.FLAG_PREV_TARGET_ATTESTER)) {
101+
vm.prev_epoch_on_chain_target_attester_hit.incr();
102+
} else {
103+
vm.prev_epoch_on_chain_target_attester_miss.incr();
104+
}
105+
106+
if (balances) |b| {
107+
if (index < b.len) total_balance += b[index];
108+
}
109+
}
110+
111+
if (balances != null) {
112+
vm.prev_epoch_on_chain_balance.set(total_balance);
113+
}
114+
}
115+
116+
test "registerValidatorStatuses records metrics" {
117+
const allocator = std.testing.allocator;
118+
try metrics.init(allocator, std.testing.io, .{});
119+
defer metrics.deinit();
120+
121+
var monitor = ValidatorMonitor.init(allocator);
122+
defer monitor.deinit();
123+
try monitor.registerLocalValidator(std.testing.io, 0);
124+
try monitor.registerLocalValidator(std.testing.io, 1);
125+
126+
const flags = [_]u8{
127+
attester_status.FLAG_PREV_SOURCE_ATTESTER |
128+
attester_status.FLAG_PREV_TARGET_ATTESTER |
129+
attester_status.FLAG_PREV_HEAD_ATTESTER,
130+
0,
131+
};
132+
const balances = [_]u64{ 32_000_000_000, 31_000_000_000 };
133+
monitor.registerValidatorStatuses(1, &flags, &balances);
134+
135+
var aw: std.Io.Writer.Allocating = .init(allocator);
136+
var list = _: {
137+
errdefer aw.deinit();
138+
try metrics.write(&aw.writer);
139+
break :_ aw.toArrayList();
140+
};
141+
defer list.deinit(allocator);
142+
143+
const expectations = [_][]const u8{
144+
"validator_monitor_prev_epoch_on_chain_source_attester_hit_total 1",
145+
"validator_monitor_prev_epoch_on_chain_source_attester_miss_total 1",
146+
"validator_monitor_prev_epoch_on_chain_target_attester_hit_total 1",
147+
"validator_monitor_prev_epoch_on_chain_target_attester_miss_total 1",
148+
"validator_monitor_prev_epoch_on_chain_head_attester_hit_total 1",
149+
"validator_monitor_prev_epoch_on_chain_head_attester_miss_total 1",
150+
"validator_monitor_prev_epoch_on_chain_balance 63000000000",
151+
};
152+
for (expectations) |expected| {
153+
std.testing.expect(std.mem.indexOf(u8, list.items, expected) != null) catch |err| {
154+
std.debug.print("expected metric not found: {s}\n", .{expected});
155+
return err;
156+
};
157+
}
158+
}
159+
160+
test "registerValidatorStatuses guards" {
161+
var monitor = ValidatorMonitor.init(std.testing.allocator);
162+
defer monitor.deinit();
163+
164+
try monitor.registerLocalValidator(0);
165+
try monitor.registerLocalValidator(2);
166+
// registering twice is a no-op
167+
try monitor.registerLocalValidator(0);
168+
try std.testing.expectEqual(@as(usize, 2), monitor.validators.count());
169+
170+
// unregistering removes; unknown index is a no-op
171+
monitor.unregisterLocalValidator(2);
172+
monitor.unregisterLocalValidator(99);
173+
try std.testing.expectEqual(@as(usize, 1), monitor.validators.count());
174+
try monitor.registerLocalValidator(2);
175+
try std.testing.expectEqual(@as(usize, 2), monitor.validators.count());
176+
177+
const flags = [_]u8{
178+
attester_status.FLAG_PREV_SOURCE_ATTESTER | attester_status.FLAG_PREV_TARGET_ATTESTER,
179+
0,
180+
attester_status.FLAG_PREV_HEAD_ATTESTER,
181+
};
182+
const balances = [_]u64{ 32_000_000_000, 31_000_000_000, 33_000_000_000 };
183+
184+
// epoch 0 is registered but has no previous epoch activity
185+
monitor.registerValidatorStatuses(std.testing.io, 0, &flags, &balances);
186+
try std.testing.expectEqual(@as(?Epoch, 0), monitor.last_registered_status_epoch);
187+
188+
monitor.registerValidatorStatuses(std.testing.io, 1, &flags, &balances);
189+
try std.testing.expectEqual(@as(?Epoch, 1), monitor.last_registered_status_epoch);
190+
191+
// same epoch twice is a no-op
192+
monitor.registerValidatorStatuses(std.testing.io, 1, &flags, &balances);
193+
try std.testing.expectEqual(@as(?Epoch, 1), monitor.last_registered_status_epoch);
194+
}

0 commit comments

Comments
 (0)