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
2 changes: 1 addition & 1 deletion pkgs/api/src/event_broadcaster.zig
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ const Checkpoint = types.Checkpoint;
const Mutex = zeam_utils.SyncMutex;

fn defaultIo() std.Io {
return std.Io.Threaded.global_single_threaded.io();
return zeam_utils.process_io.get();
}

/// Maximum size of the SSE send buffer in event_broadcaster.zig. Serialized events
Expand Down
4 changes: 2 additions & 2 deletions pkgs/api/src/lib.zig
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,8 @@ pub const MetricsError = error{
};

/// Initializes the metrics system. Must be called once at startup.
pub fn init(allocator: std.mem.Allocator) !void {
try zeam_metrics.init(allocator);
pub fn init(io: std.Io, allocator: std.mem.Allocator) !void {
try zeam_metrics.init(io, allocator);
}

/// Writes metrics to a writer (for Prometheus endpoint).
Expand Down
8 changes: 5 additions & 3 deletions pkgs/cli/src/api_server.zig
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@ const StartupStatus = enum(u8) {
/// chain is optional - if null, chain-dependent endpoints will return 503
/// (API server starts before chain initialization, so chain may not be available yet)
/// Note: Metrics are served by the separate metrics_server on a different port
pub fn startAPIServer(allocator: std.mem.Allocator, port: u16, logger_config: *LoggerConfig, chain: ?*BeamChain) !*ApiServer {
pub fn startAPIServer(io: std.Io, allocator: std.mem.Allocator, port: u16, logger_config: *LoggerConfig, chain: ?*BeamChain) !*ApiServer {
// Initialize the global event broadcaster for SSE events
// This is idempotent - safe to call even if already initialized elsewhere (e.g., node.zig)
try event_broadcaster.initGlobalBroadcaster(allocator);
Expand All @@ -60,6 +60,7 @@ pub fn startAPIServer(allocator: std.mem.Allocator, port: u16, logger_config: *L
return err;
};
ctx.* = .{
.io = io,
.allocator = allocator,
.port = port,
.logger = logger,
Expand Down Expand Up @@ -245,6 +246,7 @@ fn routeConnection(io: std.Io, connection: net.Stream, allocator: std.mem.Alloca

/// API server context
pub const ApiServer = struct {
io: std.Io,
allocator: std.mem.Allocator,
port: u16,
logger: ModuleLogger,
Expand Down Expand Up @@ -284,7 +286,7 @@ pub const ApiServer = struct {
}

fn run(self: *Self) void {
const io = std.Io.Threaded.global_single_threaded.io();
const io = self.io;
const address = net.IpAddress.parseIp4("0.0.0.0", self.port) catch |err| {
self.logger.err("failed to parse server address 0.0.0.0:{d}: {}", .{ self.port, err });
self.startup_status.store(.failed, .release);
Expand Down Expand Up @@ -621,7 +623,7 @@ pub const ApiServer = struct {

/// Handle SSE events endpoint
fn handleSSEEvents(self: *Self, stream: net.Stream) !void {
const io = std.Io.Threaded.global_single_threaded.io();
const io = self.io;
var registered = false;
errdefer if (!registered) stream.close(io);
// Set SSE headers manually by writing HTTP response
Expand Down
19 changes: 12 additions & 7 deletions pkgs/cli/src/main.zig
Original file line number Diff line number Diff line change
Expand Up @@ -301,6 +301,11 @@ pub fn main(init: std.process.Init) void {
}

fn mainInner(init: std.process.Init) !void {
// Install the process-wide blocking Io before any logging or thread spawning,
// so zeam's leaf blocking primitives use the real process Io rather than
// std's debug-only global single-threaded instance.
utils_lib.process_io.install(init.io);

var gpa = std.heap.DebugAllocator(.{}).init;
const allocator = gpa.allocator();
defer {
Expand Down Expand Up @@ -412,7 +417,7 @@ fn mainInner(init: std.process.Init) !void {
std.log.info("Successfully proved and verified all transitions", .{});
},
.beam => |beamcmd| {
api.init(allocator) catch |err| {
api.init(init.io, allocator) catch |err| {
ErrorHandler.logErrorWithOperation(err, "initialize API");
return err;
};
Expand All @@ -436,13 +441,13 @@ fn mainInner(init: std.process.Init) !void {
}

// Start metrics server (doesn't need chain reference)
metrics_server_handle = metrics_server.startMetricsServer(allocator, beamcmd.@"metrics-port", &api_logger_config) catch |err| {
metrics_server_handle = metrics_server.startMetricsServer(init.io, allocator, beamcmd.@"metrics-port", &api_logger_config) catch |err| {
ErrorHandler.logErrorWithDetails(err, "start metrics server", .{ .port = beamcmd.@"metrics-port" });
return err;
};

// Start API server early. Pass null for chain - in .beam command mode, chains are created later
api_server_handle = api_server.startAPIServer(allocator, beamcmd.@"api-port", &api_logger_config, null) catch |err| {
api_server_handle = api_server.startAPIServer(init.io, allocator, beamcmd.@"api-port", &api_logger_config, null) catch |err| {
ErrorHandler.logErrorWithDetails(err, "start API server", .{ .port = beamcmd.@"api-port" });
return err;
};
Expand Down Expand Up @@ -683,11 +688,11 @@ fn mainInner(init: std.process.Init) !void {
defer allocator.free(data_dir_3);

const db_backend = beamcmd.@"db-backend";
var db_1 = try database.Db.openBackend(allocator, logger1_config.logger(.database), data_dir_1, db_backend);
var db_1 = try database.Db.openBackend(init.io, allocator, logger1_config.logger(.database), data_dir_1, db_backend);
defer db_1.deinit();
var db_2 = try database.Db.openBackend(allocator, logger2_config.logger(.database), data_dir_2, db_backend);
var db_2 = try database.Db.openBackend(init.io, allocator, logger2_config.logger(.database), data_dir_2, db_backend);
defer db_2.deinit();
var db_3 = try database.Db.openBackend(allocator, logger3_config.logger(.database), data_dir_3, db_backend);
var db_3 = try database.Db.openBackend(init.io, allocator, logger3_config.logger(.database), data_dir_3, db_backend);
defer db_3.deinit();

// Use the same shared registry for all beam nodes
Expand Down Expand Up @@ -920,7 +925,7 @@ fn mainInner(init: std.process.Init) !void {
};

var lean_node: node.Node = undefined;
lean_node.init(allocator, &start_options) catch |err| {
lean_node.init(init.io, allocator, &start_options) catch |err| {
ErrorHandler.logErrorWithOperation(err, "initialize lean node");
return err;
};
Expand Down
5 changes: 4 additions & 1 deletion pkgs/cli/src/metrics_server.zig
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ const STARTUP_POLL_NS: u64 = 1 * std.time.ns_per_ms;
/// This is a lightweight server separate from the main API server.
/// It has no rate limiting, SSE support, or chain dependency.
pub fn startMetricsServer(
io: std.Io,
allocator: std.mem.Allocator,
port: u16,
logger_config: *LoggerConfig,
Expand All @@ -20,6 +21,7 @@ pub fn startMetricsServer(

const ctx = try allocator.create(MetricsServer);
ctx.* = .{
.io = io,
.allocator = allocator,
.port = port,
.logger = logger,
Expand Down Expand Up @@ -58,6 +60,7 @@ const StartupStatus = enum(u8) {

/// Metrics server context
pub const MetricsServer = struct {
io: std.Io,
allocator: std.mem.Allocator,
port: u16,
logger: ModuleLogger,
Expand All @@ -76,7 +79,7 @@ pub const MetricsServer = struct {
}

fn run(self: *Self) void {
const io = std.Io.Threaded.global_single_threaded.io();
const io = self.io;
const address = net.IpAddress.parseIp4("0.0.0.0", self.port) catch |err| {
self.logger.err("failed to parse server address 0.0.0.0:{d}: {}", .{ self.port, err });
self.startup_status.store(.failed, .release);
Expand Down
28 changes: 17 additions & 11 deletions pkgs/cli/src/node.zig
Original file line number Diff line number Diff line change
Expand Up @@ -168,6 +168,7 @@ pub const NodeOptions = struct {
/// A Node that encapsulates the networking, blockchain, and validator functionalities.
/// It manages the event loop, network interface, clock, and beam node.
pub const Node = struct {
io: std.Io,
loop: xev.Loop,
network: networks.EthLibp2p,
beam_node: BeamNode,
Expand Down Expand Up @@ -261,6 +262,7 @@ pub const Node = struct {
/// db directory has never been created). Set it to false when wiping a db
/// that is known to exist (genesis time mismatch case).
fn wipeAndReopenDb(
io: std.Io,
db: *database.Db,
allocator: std.mem.Allocator,
database_path: []const u8,
Expand All @@ -270,7 +272,6 @@ pub const Node = struct {
ignore_not_found: bool,
) !void {
db.deinit();
const io = std.Io.Threaded.global_single_threaded.io();
// Both backends store their working set under the same base
// directory; deleting it yields a clean slate for either engine.
const backend_dir = switch (backend) {
Expand All @@ -285,14 +286,16 @@ pub const Node = struct {
return wipe_err;
}
};
db.* = try database.Db.openBackend(allocator, logger_config.logger(.database), database_path, backend);
db.* = try database.Db.openBackend(io, allocator, logger_config.logger(.database), database_path, backend);
}

pub fn init(
self: *Self,
io: std.Io,
allocator: std.mem.Allocator,
options: *const NodeOptions,
) !void {
self.io = io;
self.allocator = allocator;
self.options = options;
self.api_server_handle = null;
Expand All @@ -301,7 +304,7 @@ pub const Node = struct {
// If path is specified load from it, otherwise use default settings
const chain_spec_owned = self.options.chain_spec != null;
const chain_spec = if (self.options.chain_spec) |path|
std.Io.Dir.cwd().readFileAlloc(std.Io.Threaded.global_single_threaded.io(), path, allocator, .limited(1024 * 1024)) catch |err| {
std.Io.Dir.cwd().readFileAlloc(io, path, allocator, .limited(1024 * 1024)) catch |err| {
self.logger.err("failed to load chain spec at '{s}': {any}", .{ path, err });
return err;
}
Expand Down Expand Up @@ -387,6 +390,7 @@ pub const Node = struct {
errdefer self.clock.deinit(allocator);

var db = try database.Db.openBackend(
io,
allocator,
options.logger_config.logger(.database),
options.database_path,
Expand All @@ -407,7 +411,7 @@ pub const Node = struct {
local_finalized_state.config.genesis_time,
chain_config.genesis.genesis_time,
});
try wipeAndReopenDb(&db, allocator, options.database_path, options.logger_config, self.logger, options.db_backend, false);
try wipeAndReopenDb(io, &db, allocator, options.database_path, options.logger_config, self.logger, options.db_backend, false);
self.logger.info("stale database wiped, starting fresh & generating genesis", .{});

local_finalized_state.deinit();
Expand All @@ -418,7 +422,7 @@ pub const Node = struct {
} else |_| {
self.logger.info("no finalized state found in db, wiping database for a clean slate", .{});
// ignore_not_found=true: db dir may not exist yet on a fresh install
try wipeAndReopenDb(&db, allocator, options.database_path, options.logger_config, self.logger, options.db_backend, true);
try wipeAndReopenDb(io, &db, allocator, options.database_path, options.logger_config, self.logger, options.db_backend, true);
self.logger.info("starting fresh & generating genesis", .{});
try self.anchor_state.genGenesisState(allocator, chain_config.genesis);
}
Expand All @@ -428,7 +432,7 @@ pub const Node = struct {
self.logger.info("checkpoint sync enabled, downloading state from: {s}", .{checkpoint_url});

// Try checkpoint sync, fall back to database/genesis on failure
if (downloadCheckpointState(allocator, checkpoint_url, self.logger)) |downloaded_state_const| {
if (downloadCheckpointState(io, allocator, checkpoint_url, self.logger)) |downloaded_state_const| {
var downloaded_state = downloaded_state_const;
// Verify state against genesis config
if (verifyCheckpointState(allocator, &downloaded_state, &chain_config.genesis, self.logger)) {
Expand Down Expand Up @@ -460,7 +464,7 @@ pub const Node = struct {
self.logger.warn("checkpoint block fetch: hashTreeRoot(BeamBlockHeader) failed: {} — skipping", .{err});
break :anchor_block_fetch;
};
downloadAndStoreCheckpointBlock(allocator, checkpoint_url, anchor_block_root, anchor_state_root, &db, self.logger);
downloadAndStoreCheckpointBlock(io, allocator, checkpoint_url, anchor_block_root, anchor_state_root, &db, self.logger);
}
} else {
self.logger.warn("skipping checkpoint sync downloaded stale/same state at slot={d}, falling back to database", .{downloaded_state.slot});
Expand Down Expand Up @@ -488,7 +492,7 @@ pub const Node = struct {
// initialization (like lean_validators_count) are captured on real
// metrics instead of being discarded by noop metrics.
if (options.metrics_enable) {
try api.init(allocator);
try api.init(io, allocator);
zeam_metrics.metrics.lean_node_start_time_seconds.set(@intCast(zeam_utils.unixTimestampSeconds()));
}

Expand Down Expand Up @@ -571,7 +575,7 @@ pub const Node = struct {

self.thread_pool = try ThreadPool.init(.{
.allocator = allocator,
.io = std.Io.Threaded.global_single_threaded.io(),
.io = io,
.thread_count = @intCast(worker_count),
});
errdefer self.thread_pool.deinit();
Expand Down Expand Up @@ -625,6 +629,7 @@ pub const Node = struct {

// Start metrics server (doesn't need chain reference)
self.metrics_server_handle = try metrics_server.startMetricsServer(
io,
allocator,
options.metrics_port,
options.logger_config,
Expand All @@ -648,6 +653,7 @@ pub const Node = struct {

// Start API server (pass chain pointer for chain-dependent endpoints)
self.api_server_handle = try api_server.startAPIServer(
io,
allocator,
options.api_port,
options.logger_config,
Expand Down Expand Up @@ -1112,13 +1118,13 @@ pub fn buildStartOptions(
/// Downloads finalized checkpoint state from the given URL and deserializes it
/// Returns the deserialized state. The caller is responsible for calling deinit on it.
fn downloadCheckpointState(
io: std.Io,
allocator: std.mem.Allocator,
url: []const u8,
logger: zeam_utils.ModuleLogger,
) !types.BeamState {
logger.info("downloading checkpoint state from: {s}", .{url});

const io = std.Io.Threaded.global_single_threaded.io();
var client = std.http.Client{
.allocator = allocator,
.io = io,
Expand Down Expand Up @@ -1254,6 +1260,7 @@ const FINALIZED_BLOCK_PATH = "/lean/v0/blocks/finalized";
/// anything — a missing anchor block is non-fatal: blocks_by_root will return
/// empty for this root until the real block arrives via reqresp or gossip.
fn downloadAndStoreCheckpointBlock(
io: std.Io,
allocator: std.mem.Allocator,
state_url: []const u8,
expected_root: types.Root,
Expand Down Expand Up @@ -1283,7 +1290,6 @@ fn downloadAndStoreCheckpointBlock(

logger.info("checkpoint block fetch: downloading anchor block from: {s}", .{block_url});

const io = std.Io.Threaded.global_single_threaded.io();
var client = std.http.Client{ .allocator = allocator, .io = io };
defer client.deinit();

Expand Down
14 changes: 8 additions & 6 deletions pkgs/database/src/database.zig
Original file line number Diff line number Diff line change
Expand Up @@ -56,11 +56,11 @@ pub const Backend = enum {
/// loud warning. Best-effort: any filesystem error here is suppressed
/// — it is a diagnostic, not a correctness barrier.
fn warnIfOtherBackendPopulated(
io: std.Io,
logger: zeam_utils.ModuleLogger,
path: []const u8,
selected: Backend,
) void {
const io = std.Io.Threaded.global_single_threaded.io();
const other: Backend = switch (selected) {
.rocksdb => .lmdb,
.lmdb => .rocksdb,
Expand Down Expand Up @@ -107,22 +107,24 @@ pub const Db = union(Backend) {
/// Open with the default backend (rocksdb). Kept for call sites
/// that don't care which engine is used (tests, legacy paths).
pub fn open(
io: std.Io,
allocator: Allocator,
logger: zeam_utils.ModuleLogger,
path: []const u8,
) !Db {
return openBackend(allocator, logger, path, .rocksdb);
return openBackend(io, allocator, logger, path, .rocksdb);
}

/// Open with an explicitly chosen backend. Used by the CLI so
/// operators can select the engine via `--db-backend`.
pub fn openBackend(
io: std.Io,
allocator: Allocator,
logger: zeam_utils.ModuleLogger,
path: []const u8,
backend: Backend,
) !Db {
warnIfOtherBackendPopulated(logger, path, backend);
warnIfOtherBackendPopulated(io, logger, path, backend);
return switch (backend) {
.rocksdb => Db{ .rocksdb = try RocksDbBackend.open(allocator, logger, path) },
.lmdb => Db{ .lmdb = try LmdbBackend.open(allocator, logger, path) },
Expand Down Expand Up @@ -504,7 +506,7 @@ fn testSaveAndLoadBlock(backend: Backend) !void {
const data_dir = try std.fmt.allocPrint(allocator, ".zig-cache/tmp/{s}", .{tmp_dir.sub_path});
defer allocator.free(data_dir);

var db = try Db.openBackend(allocator, zeam_logger_config.logger(.database_test), data_dir, backend);
var db = try Db.openBackend(std.Io.Threaded.global_single_threaded.io(), allocator, zeam_logger_config.logger(.database_test), data_dir, backend);
defer db.deinit();

const block_root = test_helpers.createDummyRoot(0xAB);
Expand Down Expand Up @@ -546,7 +548,7 @@ fn testBatchWriteAndCommit(backend: Backend) !void {
const data_dir = try std.fmt.allocPrint(allocator, ".zig-cache/tmp/{s}", .{tmp_dir.sub_path});
defer allocator.free(data_dir);

var db = try Db.openBackend(allocator, zeam_logger_config.logger(.database_test), data_dir, backend);
var db = try Db.openBackend(std.Io.Threaded.global_single_threaded.io(), allocator, zeam_logger_config.logger(.database_test), data_dir, backend);
defer db.deinit();

const block_root = test_helpers.createDummyRoot(0xAA);
Expand Down Expand Up @@ -599,7 +601,7 @@ fn testLoadLatestFinalizedStateHappyPath(backend: Backend) !void {
const data_dir = try std.fmt.allocPrint(allocator, ".zig-cache/tmp/{s}", .{tmp_dir.sub_path});
defer allocator.free(data_dir);

var db = try Db.openBackend(allocator, zeam_logger_config.logger(.database_test), data_dir, backend);
var db = try Db.openBackend(std.Io.Threaded.global_single_threaded.io(), allocator, zeam_logger_config.logger(.database_test), data_dir, backend);
defer db.deinit();

// Empty db -> no finalized slot metadata
Expand Down
2 changes: 1 addition & 1 deletion pkgs/database/src/lmdb.zig
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ pub fn Lmdb(comptime column_namespaces: []const ColumnNamespace) type {

pub fn open(allocator: Allocator, logger: zeam_utils.ModuleLogger, path: []const u8) OpenError!Self {
logger.info("initializing LMDB", .{});
const io = std.Io.Threaded.global_single_threaded.io();
const io = zeam_utils.process_io.get();

comptime {
if (column_namespaces.len == 0 or !std.mem.eql(u8, column_namespaces[0].namespace, "default")) {
Expand Down
Loading
Loading