diff --git a/v2/components/snapshot/download.zig b/v2/components/snapshot/download.zig index 82d6660f18..7b44cb10db 100644 --- a/v2/components/snapshot/download.zig +++ b/v2/components/snapshot/download.zig @@ -929,6 +929,12 @@ pub const Downloader = struct { metrics: Metrics, logger: tel.Logger("snapshot"), + // See issue #1746. + await_state: tel.BootstrapWait, + + const AWAITING_LOG_INTERVAL_NS: u64 = 10 * std.time.ns_per_s; + const AWAITING_WARN_AFTER_NS: u64 = 60 * std.time.ns_per_s; + pub fn init( gossip_to_snapshot: *SnapshotSourceRing, known_validators: KnownValidators, @@ -936,6 +942,7 @@ pub const Downloader = struct { metrics: Metrics, logger: tel.Logger("snapshot"), ) !Downloader { + const now_ns = lib.clock.monotonic(.ns); return .{ .ring = try IoUring.init(IO_URING_ENTRIES, 0), .gossip_iter = gossip_to_snapshot.get(.reader), @@ -950,6 +957,7 @@ pub const Downloader = struct { .run_result = null, .metrics = metrics, .logger = logger, + .await_state = .init(AWAITING_LOG_INTERVAL_NS, AWAITING_WARN_AFTER_NS, now_ns), }; } @@ -974,6 +982,7 @@ pub const Downloader = struct { while (true) { try self.drainGossip(); + self.maybeLogAwaitingPeers(); _ = try self.ring.submit_and_wait(0); const n = try self.ring.copy_cqes(&cqes, 0); @@ -2665,6 +2674,35 @@ pub const Downloader = struct { }; } + /// See issue #1746. + fn maybeLogAwaitingPeers(self: *Downloader) void { + if (self.download_race.phase != .idle) { + self.await_state.logReady( + self.logger, + "snapshot: found usable peer, starting download", + ); + return; + } + + const now_ns = lib.clock.monotonic(.ns); + const peers_seen = self.dedupe_map.len; + if (peers_seen == 0) { + self.await_state.logAwaiting( + now_ns, + self.logger, + "snapshot: awaiting peers from gossip ({d}s, none received)", + .{}, + ); + } else { + self.await_state.logAwaiting( + now_ns, + self.logger, + "snapshot: no usable peers so far (received={d}, active_probes={d}, {d}s)", + .{ peers_seen, self.active_probes }, + ); + } + } + /// Retires all active download connections except the one at `keep_index`. /// Used after the race reaches a terminal state (completed or failed) to /// clean up losers. Late CQEs from retired slots are ignored via gen mismatch. diff --git a/v2/lib/telemetry.zig b/v2/lib/telemetry.zig index db8706ed5c..60c19cda27 100644 --- a/v2/lib/telemetry.zig +++ b/v2/lib/telemetry.zig @@ -453,6 +453,394 @@ pub const Gauge = struct { } }; +/// Emits at most one log per `interval_ns`. Callers own the `escalated` flag +/// to promote the log level exactly once. Takes `now_ns` from a monotonic +/// clock to keep the state machine testable. +pub const ThrottledLogger = struct { + interval_ns: u64, + last_ns: u64, + escalated: bool = false, + + pub fn init(interval_ns: u64, start_ns: u64) ThrottledLogger { + return .{ .interval_ns = interval_ns, .last_ns = start_ns }; + } + + pub fn tick(self: *ThrottledLogger, now_ns: u64) bool { + // Saturating subtract guards against a hypothetical clock regression. + if (now_ns -| self.last_ns >= self.interval_ns) { + self.last_ns = now_ns; + return true; + } + return false; + } +}; + +test "ThrottledLogger: does not fire until interval has elapsed" { + var t: ThrottledLogger = .init(1_000, 0); + try std.testing.expect(!t.tick(0)); + try std.testing.expect(!t.tick(1)); + try std.testing.expect(!t.tick(999)); + try std.testing.expect(t.tick(1_000)); +} + +test "ThrottledLogger: subsequent ticks respect the interval" { + var t: ThrottledLogger = .init(1_000, 0); + try std.testing.expect(t.tick(1_000)); + try std.testing.expect(!t.tick(1_500)); + try std.testing.expect(t.tick(2_001)); + try std.testing.expect(!t.tick(2_500)); + try std.testing.expect(t.tick(3_500)); +} + +test "ThrottledLogger: escalated flag is caller-owned and starts false" { + var t: ThrottledLogger = .init(1_000, 0); + try std.testing.expect(!t.escalated); + t.escalated = true; + try std.testing.expect(t.escalated); + try std.testing.expect(t.tick(2_000)); +} + +test "ThrottledLogger: start_ns suppresses early ticks" { + var t: ThrottledLogger = .init(1_000, 50_000); + try std.testing.expect(!t.tick(50_000)); + try std.testing.expect(!t.tick(50_500)); + try std.testing.expect(t.tick(51_000)); +} + +/// Throttled awaiting-log + one-shot warn escalation + one-shot ready log, +/// used by v2 bootstrap services blocked on an upstream signal. See #1746. +pub const BootstrapWait = struct { + gate: ThrottledLogger, + start_ns: u64, + warn_after_ns: u64, + logged_awaiting: bool = false, + logged_ready: bool = false, + + pub const Action = union(enum) { + none, + info: u64, + /// Returned at most once; subsequent overdue ticks return `.info`. + warn: u64, + }; + + pub fn init(interval_ns: u64, warn_after_ns: u64, start_ns: u64) BootstrapWait { + return .{ + .gate = .init(interval_ns, start_ns), + .start_ns = start_ns, + .warn_after_ns = warn_after_ns, + }; + } + + pub fn tick(self: *BootstrapWait, now_ns: u64) Action { + if (!self.gate.tick(now_ns)) return .none; + self.logged_awaiting = true; + const elapsed_ns = now_ns -| self.start_ns; + const elapsed_s = elapsed_ns / std.time.ns_per_s; + if (!self.gate.escalated and elapsed_ns >= self.warn_after_ns) { + self.gate.escalated = true; + return .{ .warn = elapsed_s }; + } + return .{ .info = elapsed_s }; + } + + /// Returns true exactly once, and only if a prior `tick` fired. + pub fn markReady(self: *BootstrapWait) bool { + if (self.logged_awaiting and !self.logged_ready) { + self.logged_ready = true; + return true; + } + return false; + } + + /// `tick` + emit. `fmt` must end with `({d}s)`; elapsed_s is appended last. + pub fn logAwaiting( + self: *BootstrapWait, + now_ns: u64, + logger: anytype, + comptime fmt: []const u8, + args: anytype, + ) void { + switch (self.tick(now_ns)) { + .none => {}, + .info => |s| logger.info().logf(fmt, args ++ .{s}), + .warn => |s| logger.warn().logf(fmt, args ++ .{s}), + } + } + + pub fn logReady( + self: *BootstrapWait, + logger: anytype, + comptime msg: []const u8, + ) void { + if (self.markReady()) logger.info().log(msg); + } +}; + +/// `View.getBufferBlocking(runner)` with bootstrap-awaiting logs woven in. +/// Skips the ready log on empty-slice EOF (see `v2/lib/ipc/ring.zig:60-64`). +pub fn waitForBufferWithAwaitingLog( + view: anytype, + runner: anytype, + wait: *BootstrapWait, + logger: anytype, + comptime awaiting_fmt: []const u8, + comptime ready_msg: []const u8, +) !@typeInfo(@TypeOf(view.getBuffer())).optional.child { + if (view.getBuffer()) |buf| { + if (buf.len != 0) wait.logReady(logger, ready_msg); + return buf; + } + const buf = while (true) { + wait.logAwaiting(clock.monotonic(.ns), logger, awaiting_fmt, .{}); + try runner.activity.signalIdleSpinning(); + if (view.getBuffer()) |b| break b; + }; + try runner.activity.signalActive(); + if (buf.len != 0) wait.logReady(logger, ready_msg); + return buf; +} + +test "BootstrapWait: no tick fires before interval elapses" { + var w: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, 0); + try std.testing.expectEqual(BootstrapWait.Action.none, w.tick(0)); + try std.testing.expectEqual(BootstrapWait.Action.none, w.tick(9 * std.time.ns_per_s)); + try std.testing.expect(!w.logged_awaiting); +} + +test "BootstrapWait: first tick after interval returns .info" { + var w: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, 0); + const action = w.tick(10 * std.time.ns_per_s); + try std.testing.expectEqual(@as(u64, 10), action.info); + try std.testing.expect(w.logged_awaiting); + try std.testing.expect(!w.gate.escalated); +} + +test "BootstrapWait: escalates to .warn exactly once at warn_after_ns" { + var w: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, 0); + try std.testing.expectEqual(BootstrapWait.Action{ .info = 10 }, w.tick(10 * std.time.ns_per_s)); + try std.testing.expectEqual(BootstrapWait.Action{ .info = 30 }, w.tick(30 * std.time.ns_per_s)); + try std.testing.expectEqual(BootstrapWait.Action{ .warn = 60 }, w.tick(60 * std.time.ns_per_s)); + try std.testing.expect(w.gate.escalated); + try std.testing.expectEqual( + BootstrapWait.Action{ .info = 90 }, + w.tick(90 * std.time.ns_per_s), + ); + try std.testing.expectEqual( + BootstrapWait.Action{ .info = 120 }, + w.tick(120 * std.time.ns_per_s), + ); +} + +test "BootstrapWait: throttle applies between ticks" { + var w: BootstrapWait = .init(10 * std.time.ns_per_s, 600 * std.time.ns_per_s, 0); + try std.testing.expectEqual(BootstrapWait.Action{ .info = 10 }, w.tick(10 * std.time.ns_per_s)); + try std.testing.expectEqual(BootstrapWait.Action.none, w.tick(15 * std.time.ns_per_s)); + try std.testing.expectEqual(BootstrapWait.Action.none, w.tick(19 * std.time.ns_per_s)); + try std.testing.expectEqual(BootstrapWait.Action{ .info = 20 }, w.tick(20 * std.time.ns_per_s)); +} + +test "BootstrapWait: escalation deadline before first interval still triggers .warn on first fire" { + var w: BootstrapWait = .init(10 * std.time.ns_per_s, 1 * std.time.ns_per_s, 0); + try std.testing.expectEqual(BootstrapWait.Action{ .warn = 10 }, w.tick(10 * std.time.ns_per_s)); + try std.testing.expect(w.gate.escalated); +} + +test "BootstrapWait: markReady returns false when awaiting never logged" { + var w: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, 0); + try std.testing.expect(!w.markReady()); + _ = w.tick(5 * std.time.ns_per_s); + try std.testing.expect(!w.markReady()); +} + +test "BootstrapWait: markReady returns true once after tick fired" { + var w: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, 0); + _ = w.tick(10 * std.time.ns_per_s); + try std.testing.expect(w.markReady()); + try std.testing.expect(!w.markReady()); + try std.testing.expect(!w.markReady()); +} + +test "BootstrapWait: markReady still gated by logged_awaiting after later ticks" { + const start_ns = 100 * std.time.ns_per_s; + var w: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, start_ns); + _ = w.tick(start_ns); + try std.testing.expect(!w.markReady()); + _ = w.tick(start_ns + 11 * std.time.ns_per_s); + try std.testing.expect(w.markReady()); + try std.testing.expect(!w.markReady()); +} + +test "BootstrapWait: saturating subtract survives clock regression" { + const start_ns = 100 * std.time.ns_per_s; + var w: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, start_ns); + try std.testing.expectEqual(BootstrapWait.Action.none, w.tick(50 * std.time.ns_per_s)); + try std.testing.expectEqual(BootstrapWait.Action.none, w.tick(start_ns - 1)); + try std.testing.expectEqual( + BootstrapWait.Action{ .info = 10 }, + w.tick(start_ns + 10 * std.time.ns_per_s), + ); +} + +test "BootstrapWait: logAwaiting dispatches all three branches through a noop logger" { + const logger = Logger("test").noop; + var w: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, 0); + + w.logAwaiting(5 * std.time.ns_per_s, logger, "test: waiting ({d}s)", .{}); + try std.testing.expect(!w.logged_awaiting); + + w.logAwaiting(10 * std.time.ns_per_s, logger, "test: waiting ({d}s)", .{}); + try std.testing.expect(w.logged_awaiting); + try std.testing.expect(!w.gate.escalated); + + w.logAwaiting(60 * std.time.ns_per_s, logger, "test: waiting ({d}s)", .{}); + try std.testing.expect(w.gate.escalated); + + w.logAwaiting(90 * std.time.ns_per_s, logger, "test: waiting ({d}s)", .{}); + try std.testing.expect(w.gate.escalated); +} + +test "BootstrapWait: logAwaiting appends elapsed_s to arbitrary leading args" { + const logger = Logger("test").noop; + var w: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, 0); + const peers_seen: usize = 3; + const active_probes: u8 = 2; + w.logAwaiting( + 10 * std.time.ns_per_s, + logger, + "test: no usable peers (received={d}, active_probes={d}, {d}s)", + .{ peers_seen, active_probes }, + ); + try std.testing.expect(w.logged_awaiting); +} + +test "BootstrapWait: logReady is a no-op when awaiting was never logged" { + const logger = Logger("test").noop; + var w: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, 0); + w.logReady(logger, "test: ready"); + try std.testing.expect(!w.logged_ready); +} + +test "BootstrapWait: logReady emits exactly once after a prior awaiting log" { + const logger = Logger("test").noop; + var w: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, 0); + w.logAwaiting(10 * std.time.ns_per_s, logger, "test: waiting ({d}s)", .{}); + w.logReady(logger, "test: ready"); + try std.testing.expect(w.logged_ready); + w.logReady(logger, "test: ready"); + try std.testing.expect(w.logged_ready); +} + +const TestMocks = struct { + const View = struct { + script: []const ?[]const u8, + cursor: usize = 0, + pub fn getBuffer(self: *View) ?[]const u8 { + const i = self.cursor; + self.cursor += 1; + return if (i < self.script.len) self.script[i] else null; + } + }; + const Activity = struct { + idle_calls: u32 = 0, + active_calls: u32 = 0, + fail_after: ?u32 = null, + pub fn signalIdleSpinning(self: *Activity) !void { + if (self.fail_after) |n| if (self.idle_calls >= n) return error.Canceled; + self.idle_calls += 1; + } + pub fn signalActive(self: *Activity) !void { + self.active_calls += 1; + } + }; + const Runner = struct { + activity: *Activity, + }; +}; + +test "waitForBufferWithAwaitingLog: fast path returns immediately without signaling" { + var view: TestMocks.View = .{ .script = &.{"hello"} }; + var activity: TestMocks.Activity = .{}; + const runner: TestMocks.Runner = .{ .activity = &activity }; + var wait: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, 0); + const logger = Logger("test").noop; + + const buf = try waitForBufferWithAwaitingLog( + &view, + runner, + &wait, + logger, + "test: awaiting ({d}s)", + "test: ready", + ); + + try std.testing.expectEqualStrings("hello", buf); + try std.testing.expectEqual(@as(u32, 0), activity.idle_calls); + try std.testing.expectEqual(@as(u32, 0), activity.active_calls); + try std.testing.expect(!wait.logged_awaiting); + try std.testing.expect(!wait.logged_ready); +} + +test "waitForBufferWithAwaitingLog: slow path signals idle then active, emits ready log" { + var view: TestMocks.View = .{ .script = &.{ null, null, null, "delayed" } }; + var activity: TestMocks.Activity = .{}; + const runner: TestMocks.Runner = .{ .activity = &activity }; + var wait: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, 0); + const logger = Logger("test").noop; + + const buf = try waitForBufferWithAwaitingLog( + &view, + runner, + &wait, + logger, + "test: awaiting ({d}s)", + "test: ready", + ); + + try std.testing.expectEqualStrings("delayed", buf); + try std.testing.expect(activity.idle_calls >= 1); + try std.testing.expectEqual(@as(u32, 1), activity.active_calls); +} + +test "waitForBufferWithAwaitingLog: empty-slice EOF does not trigger ready log" { + const empty: []const u8 = &.{}; + var view: TestMocks.View = .{ .script = &.{empty} }; + var activity: TestMocks.Activity = .{}; + const runner: TestMocks.Runner = .{ .activity = &activity }; + var wait: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, 0); + wait.logAwaiting(10 * std.time.ns_per_s, Logger("test").noop, "test: awaiting ({d}s)", .{}); + + const buf = try waitForBufferWithAwaitingLog( + &view, + runner, + &wait, + Logger("test").noop, + "test: awaiting ({d}s)", + "test: ready", + ); + + try std.testing.expectEqual(@as(usize, 0), buf.len); + try std.testing.expect(!wait.logged_ready); +} + +test "waitForBufferWithAwaitingLog: signalIdleSpinning error propagates" { + var view: TestMocks.View = .{ .script = &.{null} }; + var activity: TestMocks.Activity = .{ .fail_after = 0 }; + const runner: TestMocks.Runner = .{ .activity = &activity }; + var wait: BootstrapWait = .init(10 * std.time.ns_per_s, 60 * std.time.ns_per_s, 0); + const logger = Logger("test").noop; + + const result = waitForBufferWithAwaitingLog( + &view, + runner, + &wait, + logger, + "test: awaiting ({d}s)", + "test: ready", + ); + + try std.testing.expectError(error.Canceled, result); +} + /// Can be used as a counter or a gauge. pub fn Variant(comptime V: type) type { return struct { diff --git a/v2/services/accounts_db.zig b/v2/services/accounts_db.zig index 201f24aded..448452a7fa 100644 --- a/v2/services/accounts_db.zig +++ b/v2/services/accounts_db.zig @@ -20,6 +20,11 @@ pub const std_options = start.options; pub const ReadOnly = services.accounts_db.ReadOnly; pub const ReadWrite = services.accounts_db.ReadWrite; +const ServiceLogger = lib.telemetry.Logger("main"); + +const SNAPSHOT_WAIT_INTERVAL_NS: u64 = 10 * std.time.ns_per_s; +const SNAPSHOT_WAIT_WARN_AFTER_NS: u64 = 5 * 60 * std.time.ns_per_s; + pub fn serviceMain(runner: lib.runner.Connection, _: ReadOnly, rw: ReadWrite) !noreturn { const logger = rw.tel.acquireLogger(@tagName(name), "main"); rw.tel.signalReady(); @@ -30,6 +35,8 @@ pub fn serviceMain(runner: lib.runner.Connection, _: ReadOnly, rw: ReadWrite) !n const Global = struct { var fba_memory: [32 * 1024 * 1024]u8 = undefined; var rooted: Rooted = undefined; + // Held by pointer on `SnapshotBufReader` to survive pass-by-value copies. + var wait_state: lib.telemetry.BootstrapWait = undefined; }; const rooted = &Global.rooted; @@ -50,20 +57,34 @@ pub fn serviceMain(runner: lib.runner.Connection, _: ReadOnly, rw: ReadWrite) !n if (rooted.table.count() == 0) { logger.info().log("no existing rooted db. reading from snapshot"); + const now_ns = lib.clock.monotonic(.ns); + Global.wait_state = .init( + SNAPSHOT_WAIT_INTERVAL_NS, + SNAPSHOT_WAIT_WARN_AFTER_NS, + now_ns, + ); + const SnapshotDataRingReader = @TypeOf(in); const SnapshotBufReader = struct { in_: *SnapshotDataRingReader, runner_: lib.runner.Connection, completion_: *std.atomic.Value(f64), + wait_: *lib.telemetry.BootstrapWait, + logger_: ServiceLogger, pub fn percentCompleted(self: @This()) f64 { return self.completion_.load(.monotonic); } pub fn getBuffer(self: @This()) []const u8 { - return self.in_.getBufferBlocking(self.runner_) catch |err| switch (err) { - error.Canceled => return &.{}, // cancel -> EOF - }; + return lib.telemetry.waitForBufferWithAwaitingLog( + self.in_, + self.runner_, + self.wait_, + self.logger_, + "accounts_db: awaiting snapshot bytes from snapshot service ({d}s)", + "accounts_db: receiving snapshot bytes", + ) catch return &.{}; } pub fn advance(self: @This(), n: usize) void { @@ -76,6 +97,8 @@ pub fn serviceMain(runner: lib.runner.Connection, _: ReadOnly, rw: ReadWrite) !n .in_ = &in, .runner_ = runner, .completion_ = &rw.ready_snapshot_in.completion, + .wait_ = &Global.wait_state, + .logger_ = logger, }); logger.info().log("reading snapshot accounts"); diff --git a/v2/services/gossip.zig b/v2/services/gossip.zig index 9f30e31192..7466714c6d 100644 --- a/v2/services/gossip.zig +++ b/v2/services/gossip.zig @@ -112,6 +112,16 @@ pub fn serviceMain(runner: lib.runner.Connection, ro: ReadOnly, rw: ReadWrite) ! }); var it = rw.net_pair.recv.get(.reader); + + // See issue #1746. + const bootstrap_start_ns = lib.clock.monotonic(.ns); + var wait_state: lib.telemetry.BootstrapWait = .init( + 10 * std.time.ns_per_s, + 60 * std.time.ns_per_s, + bootstrap_start_ns, + ); + var first_packet_received = false; + while (true) { now = lib.clock.wallclock(.ms); try gossip_node.poll(.from(logger), now); @@ -122,10 +132,33 @@ pub fn serviceMain(runner: lib.runner.Connection, ro: ReadOnly, rw: ReadWrite) ! // For now this should work fine, but in theory there's a very slim chance // of a race condition (it should be basically impossible to manifest // in the one black-box test that currently exists for this). + + if (!first_packet_received) { + wait_state.logAwaiting( + lib.clock.monotonic(.ns), + logger, + "gossip has received no packets from cluster ({d}s)", + .{}, + ); + } + try runner.activity.signalIdleSpinning(); continue; }; try runner.activity.signalActive(); + + if (!first_packet_received) { + first_packet_received = true; + // Milestone log fires unconditionally; ignore markReady's gate. + _ = wait_state.markReady(); + const elapsed_s = + (lib.clock.monotonic(.ns) -| bootstrap_start_ns) / std.time.ns_per_s; + logger.info().logf( + "gossip received first packets from cluster ({d}s)", + .{elapsed_s}, + ); + } + gossip_node.processPacket(.from(logger), now, packet); it.markUsed(); } diff --git a/v2/services/replay.zig b/v2/services/replay.zig index 55c571e9bd..6d1f5eaca1 100644 --- a/v2/services/replay.zig +++ b/v2/services/replay.zig @@ -275,6 +275,27 @@ pub fn serviceMain(runner: lib.runner.Connection, _: ReadOnly, rw: ReadWrite) !n } } +// Warn threshold sits past the realistic tail of mainnet snapshot loading +// so operators aren't paged during normal startup. See issue #1746. +const RUNTIME_METADATA_WAIT_INTERVAL_NS: u64 = 10 * std.time.ns_per_s; +const RUNTIME_METADATA_WAIT_WARN_AFTER_NS: u64 = 15 * 60 * std.time.ns_per_s; + +fn waitForBlockhashes( + blockhashes_in: anytype, + runner: lib.runner.Connection, + logger: tel.Logger("main"), + wait_state: *tel.BootstrapWait, +) ![]const Hash { + return tel.waitForBufferWithAwaitingLog( + blockhashes_in, + runner, + wait_state, + logger, + "replay: awaiting runtime metadata from accounts_db ({d}s)", + "replay: receiving runtime metadata", + ); +} + /// Reads all the RuntimeMetadata provided by accountsdb from the snapshot or /// its internal state. This bootstraps replay with information about its /// starting root slot, and some older info like the history of blockhashes. @@ -299,6 +320,13 @@ fn bootstrap( exec_states: *BlockExecStates, blockhash_states: *BlockHashStates, ) !void { + const now_ns = lib.clock.monotonic(.ns); + var wait_state: tel.BootstrapWait = .init( + RUNTIME_METADATA_WAIT_INTERVAL_NS, + RUNTIME_METADATA_WAIT_WARN_AFTER_NS, + now_ns, + ); + var num_hashes: usize = 0; // Drain the blockhash queue into the block tree. accountsdb writes into // this ring blocks waiting for the reader (us). @@ -307,7 +335,7 @@ fn bootstrap( defer blockhashes_in.close(); var last_block: ?api.BlockRef = null; while (true) { - const hashes = try blockhashes_in.getBufferBlocking(runner); + const hashes = try waitForBlockhashes(&blockhashes_in, runner, logger, &wait_state); if (hashes.len == 0) break; // blockhashes_out closed their end for (hashes) |*hash| { const block = try block_pool.createId();