diff --git a/pkgs/node/src/chain.zig b/pkgs/node/src/chain.zig index 52d914853..d2867b5c7 100644 --- a/pkgs/node/src/chain.zig +++ b/pkgs/node/src/chain.zig @@ -817,7 +817,22 @@ pub const BeamChain = struct { // import block assuming it is gossip validated or synced // this onBlock corresponds to spec's forkchoice's onblock with some functionality split between this and - // our implemented forkchoice's onblock. this is to parallelize "apply transition" with other verifications + // our implemented forkchoice's onblock. + // + // Parallelization plan (tracked in https://github.com/blockblaz/zeam/issues/719): + // + // Phase 1 – Parallel validations (this PR scaffolds; full impl in follow-up): + // a) verifySignatures — CPU-bound XMSS sig verification, independent of STF + // b) apply_transition — state transition, independent of sig verification result + // Both are dispatched to spawned threads; forkchoice import is gated on both completing + // successfully. + // + // Phase 2 – Parallel post-state compute: + // Independent segments of state computation identified and parallelised. + // + // Phase 3 – Execution / proof verification (future): + // Will be parallelised once Phase 1/2 are stable. + // // Returns a list of missing block roots that need to be fetched from the network pub fn onBlock(self: *Self, signedBlock: types.SignedBlock, blockInfo: CachedProcessedBlockInfo) ![]types.Root { const onblock_timer = zeam_metrics.chain_onblock_duration_seconds.start(); @@ -843,16 +858,80 @@ pub const BeamChain = struct { // If anything below fails, deinit interior first (LIFO: deinit runs before destroy above). errdefer cpost_state.deinit(); - // 2. verify XMSS signatures (independent step; placed before STF for now, parallelizable later) - // Use public key cache to avoid repeated SSZ deserialization of validator public keys - try stf.verifySignatures(self.allocator, pre_state, &signedBlock, &self.public_key_cache); + // 2+3. verify XMSS signatures (task A) and apply state transition (task B) in parallel (issue #719) + + // Task A gets its own temporary arena so its allocations are freed after it completes, + // keeping peak memory equivalent to the sequential baseline (avoids OOM on CI). + // Task B uses the chain allocator directly since its output (cpost_state) must outlive the task. + var sig_arena = std.heap.ArenaAllocator.init(self.allocator); + defer sig_arena.deinit(); + + // Task A: verify XMSS signatures + const SigVerifyTask = struct { + allocator: Allocator, + pre_state: *const types.BeamState, + signed_block: *const types.SignedBlock, + pubkey_cache: *xmss.PublicKeyCache, + err: ?anyerror = null, + + fn run(task: *@This()) void { + stf.verifySignatures(task.allocator, task.pre_state, task.signed_block, task.pubkey_cache) catch |e| { + task.err = e; + }; + } + }; + var sig_task = SigVerifyTask{ + .allocator = sig_arena.allocator(), + .pre_state = pre_state, + .signed_block = &signedBlock, + .pubkey_cache = &self.public_key_cache, + }; + + // Task B: apply state transition + const StfTask = struct { + allocator: Allocator, + post_state: *types.BeamState, + block: types.BeamBlock, + stf_logger: zeam_utils.ModuleLogger, + root_to_slot_cache: *types.RootToSlotCache, + err: ?anyerror = null, + + fn run(task: *@This()) void { + stf.apply_transition(task.allocator, task.post_state, task.block, .{ + .logger = task.stf_logger, + .validSignatures = true, + .rootToSlotCache = task.root_to_slot_cache, + }) catch |e| { + task.err = e; + }; + } + }; + var stf_task = StfTask{ + .allocator = self.allocator, + .post_state = cpost_state, + .block = block, + .stf_logger = self.stf_logger, + .root_to_slot_cache = &self.root_to_slot_cache, + }; + + // Spawn threads; fall back to inline execution if spawn fails. + const sig_thread = std.Thread.spawn(.{}, SigVerifyTask.run, .{&sig_task}) catch null_thread: { + SigVerifyTask.run(&sig_task); + break :null_thread null; + }; + const stf_thread = std.Thread.spawn(.{}, StfTask.run, .{&stf_task}) catch null_thread: { + StfTask.run(&stf_task); + break :null_thread null; + }; + + // Join successfully spawned threads (barrier). + if (sig_thread) |t| t.join(); + if (stf_thread) |t| t.join(); + + // Propagate any errors from parallel tasks. + if (sig_task.err) |e| return e; + if (stf_task.err) |e| return e; - // 3. apply state transition assuming signatures are valid (STF does not re-verify) - try stf.apply_transition(self.allocator, cpost_state, block, .{ - .logger = self.stf_logger, - .validSignatures = true, - .rootToSlotCache = &self.root_to_slot_cache, - }); break :computedstate cpost_state; }; // If post_state was freshly allocated above and a later step errors (e.g. forkChoice.onBlock,