From 6271aa4b023371f1e7c34a3dffd6d329a88ef1b0 Mon Sep 17 00:00:00 2001 From: xylophonez Date: Fri, 21 Aug 2026 23:12:19 +0100 Subject: [PATCH 1/5] fix(arweave-scheduler): recover from chain reorganizations The scheduler validates each block's previous_block against the last indexed indep_hash. If it ever indexes a block the network later orphans, every canonical successor fails with 'does not extend the indexed chain' (409) forever - the frontier wedges permanently and no process can advance to tip. Observed twice on real deployments; purge+resync only re-arms it. Add bounded reorg recovery: on the conflict, walk the canonical chain down (bounded by arweave-scheduler-reorg-depth, default 50; never below 'from'), comparing network indep_hash to the locally indexed block, replacing orphaned cached blocks above the fork point, rewinding the global record to the fork, and re-indexing the canonical tail. Attempted once per pass; fails closed on no match or a fetch error. Adds reorg_fork_point_test, reorg_recovery_integration_test, and a live reorg_canonical_fetch_shape_test_. --- .../process/dev_arweave_scheduler_sync.erl | 238 +++++++++++++++++- 1 file changed, 235 insertions(+), 3 deletions(-) diff --git a/src/preloaded/process/dev_arweave_scheduler_sync.erl b/src/preloaded/process/dev_arweave_scheduler_sync.erl index b11719d4f..365acc067 100644 --- a/src/preloaded/process/dev_arweave_scheduler_sync.erl +++ b/src/preloaded/process/dev_arweave_scheduler_sync.erl @@ -18,6 +18,7 @@ -define(DEFAULT_FETCH_ATTEMPTS, 20). -define(DEFAULT_FETCH_RETRY_MS, 250). -define(DEFAULT_SYNC_INTERVAL_MS, 1000). +-define(DEFAULT_REORG_DEPTH, 50). -define(PROCESS_IDLE_MS, 1000). %% @doc Bring the global target index to the confirmed chain frontier. The @@ -129,9 +130,12 @@ initial_state({ok, State}, From) -> initial_state({error, not_found}, From) -> initial_state(not_found, From); initial_state(Error, _From) -> Error. -sync_blocks(State = #{ <<"to">> := To }, Upper, _Opts) when To >= Upper -> +sync_blocks(State, Upper, Opts) -> + sync_blocks(State, Upper, Opts, _RecoveryAllowed = true). + +sync_blocks(State = #{ <<"to">> := To }, Upper, _Opts, _Recover) when To >= Upper -> {ok, State}; -sync_blocks(State = #{ <<"to">> := To }, Upper, Opts) -> +sync_blocks(State = #{ <<"to">> := To }, Upper, Opts, Recover) -> BatchEnd = min( Upper, @@ -147,6 +151,22 @@ sync_blocks(State = #{ <<"to">> := To }, Upper, Opts) -> ) ) ), + case sync_batch(State, BatchEnd, Opts) of + {ok, NewState} -> + sync_blocks(NewState, Upper, Opts, Recover); + {error, #{ + <<"status">> := 409, + <<"reason">> := + <<"Arweave block does not extend the indexed chain.">> + }} = Conflict when Recover -> + case recover_from_reorg(State, Opts) of + {ok, Rewound} -> sync_blocks(Rewound, Upper, Opts, false); + not_recovered -> Conflict + end; + Error -> Error + end. + +sync_batch(State = #{ <<"to">> := To }, BatchEnd, Opts) -> Heights = lists:seq(To + 1, BatchEnd), maybe {ok, Blocks} ?= fetch_blocks(Heights, Opts), @@ -155,9 +175,113 @@ sync_blocks(State = #{ <<"to">> := To }, Upper, Opts) -> ok ?= write_blocks(Blocks, Opts), {ok, _} ?= dev_arweave_scheduler_cache:write_global(NewState, Opts), - sync_blocks(NewState, Upper, Opts) + {ok, NewState} end. +%% @doc The indexed tip no longer matches the canonical chain: the node +%% indexed a block that the network later orphaned in a reorganization, so +%% every canonical successor fails the previous-hash check forever. Walk the +%% canonical chain downward -- a bounded number of heights, fetched from the +%% network rather than the local cache -- until a canonical block hash matches +%% the indexed block at the same height. Replace each orphaned cached block +%% with its canonical counterpart on the way down so the forward resync +%% cannot re-index the orphan, rewind the global record to the fork point, +%% and let the ordinary forward path re-index the tail. If no indexed block +%% within the bound matches the canonical chain, leave the record untouched +%% and surface the original conflict. +recover_from_reorg(State, Opts) -> + recover_from_reorg( + State, + Opts, + fun(Height) -> fetch_block_remote(Height, Opts) end + ). + +recover_from_reorg(State = #{ <<"to">> := To, <<"from">> := From }, Opts, FetchBlock) -> + Floor = max(From, To - reorg_depth(Opts)), + Result = + fork_point( + To, + Floor, + fun(Height) -> + maybe + {ok, Canonical} ?= FetchBlock(Height), + CanonicalHash = + hb_maps:get( + <<"indep_hash">>, Canonical, not_found, Opts + ), + true ?= is_binary(CanonicalHash), + {ok, CanonicalHash, Canonical} + else + _ -> error + end + end, + fun(Height) -> + case dev_arweave_scheduler_cache:read_block(Height, Opts) of + {ok, Indexed} -> + case + hb_maps:get( + <<"indep_hash">>, Indexed, not_found, Opts + ) + of + Hash when is_binary(Hash) -> {ok, Hash}; + _ -> not_found + end; + _ -> not_found + end + end, + fun(Height, Canonical) -> + dev_arweave_scheduler_cache:write_block( + Height, Canonical, Opts + ) + end + ), + case Result of + {ok, ForkHeight, ForkHash} -> + Rewound = + State#{ + <<"to">> => ForkHeight, + <<"block-hash">> => ForkHash + }, + case dev_arweave_scheduler_cache:write_global(Rewound, Opts) of + {ok, _} -> {ok, Rewound}; + _ -> not_recovered + end; + not_recovered -> not_recovered + end. + +%% @doc Find the deepest height, no lower than `Floor', whose canonical +%% network block matches the locally indexed block. Orphaned cached blocks +%% encountered above the fork point are replaced with their canonical +%% counterparts as the walk descends. +fork_point(Height, Floor, _Canonical, _Indexed, _Replace) when Height < Floor -> + not_recovered; +fork_point(Height, Floor, Canonical, Indexed, Replace) -> + case Canonical(Height) of + {ok, CanonicalHash, CanonicalBlock} -> + case Indexed(Height) of + {ok, CanonicalHash} -> {ok, Height, CanonicalHash}; + _ -> + Replace(Height, CanonicalBlock), + fork_point(Height - 1, Floor, Canonical, Indexed, Replace) + end; + error -> not_recovered + end. + +%% @doc Return the bounded number of blocks a reorganization recovery may +%% rewind. The bound keeps recovery finite and fails the sync closed when a +%% divergence is deeper than any plausible reorganization. +reorg_depth(Opts) -> + max( + 1, + hb_util:int( + hb_opts:get( + arweave_scheduler_reorg_depth, + ?DEFAULT_REORG_DEPTH, + Opts + ) + ) + ). + %% @doc Fetch one bounded batch of canonical blocks concurrently. fetch_blocks(Heights, Opts) -> Results = @@ -1117,6 +1241,114 @@ stop_test_runner(Name) -> hb_name:unregister(Name) end. +reorg_canonical_fetch_shape_test_() -> + {timeout, 60, fun() -> + Store = hb_test_utils:test_store( + hb_store_volatile, <<"ar-sched-fetchshape">>), + ok = hb_store:start(Store), + Opts = #{ <<"store">> => [Store], <<"scheduler-store">> => [Store], + <<"gateway">> => <<"https://arweave.net">>, + <<"arweave-scheduler-fetch-attempts">> => 2, + <<"priv-wallet">> => ar_wallet:new() }, + %% Recovery reads indep_hash and previous_block off blocks returned by + %% fetch_block_remote; assert two adjacent canonical blocks carry a + %% binary indep_hash and chain correctly, so the recovery closure's + %% assumption holds against live data. Skip cleanly when offline. + case {fetch_block_remote(1984251, Opts), + fetch_block_remote(1984252, Opts)} of + {{ok, B1}, {ok, B2}} -> + H1 = hb_maps:get(<<"indep_hash">>, B1, not_found, Opts), + H2 = hb_maps:get(<<"indep_hash">>, B2, not_found, Opts), + Prev2 = hb_maps:get(<<"previous_block">>, B2, not_found, Opts), + ?assert(is_binary(H1)), + ?assert(is_binary(H2)), + ?assertEqual(H1, Prev2); + _ -> + ?debugMsg("skipped: canonical block fetch unavailable"), + ok + end, + ok = hb_store:stop(Store) + end}. + +reorg_recovery_integration_test() -> + Store = hb_test_utils:test_store(hb_store_volatile, <<"ar-sched-reorg">>), + ok = hb_store:start(Store), + Opts = #{ <<"store">> => [Store], <<"scheduler-store">> => [Store] }, + Mk = fun(H, Hash, Prev) -> + #{ <<"height">> => H, <<"indep_hash">> => Hash, + <<"previous_block">> => Prev } + end, + %% Canonical chain 100..103. The node indexed canonical 100/101 but an + %% ORPHAN at 102, then wedged trying to extend with canonical 103. + Canon = fun(H) -> <<"canon-", (hb_util:bin(H))/binary>> end, + {ok, _} = dev_arweave_scheduler_cache:write_block( + 100, Mk(100, Canon(100), <<"canon-99">>), Opts), + {ok, _} = dev_arweave_scheduler_cache:write_block( + 101, Mk(101, Canon(101), Canon(100)), Opts), + {ok, _} = dev_arweave_scheduler_cache:write_block( + 102, Mk(102, <<"orphan-102">>, Canon(101)), Opts), + State0 = #{ <<"from">> => 100, <<"to">> => 102, + <<"block-hash">> => <<"orphan-102">> }, + {ok, _} = dev_arweave_scheduler_cache:write_global(State0, Opts), + %% Canonical fetcher returns the true chain (as the network would). + Fetch = fun(H) -> {ok, Mk(H, Canon(H), + case H of 100 -> <<"canon-99">>; _ -> Canon(H - 1) end)} end, + {ok, Rewound} = recover_from_reorg(State0, Opts, Fetch), + %% Fork point is 101 (deepest indexed block matching canonical). + ?assertEqual(101, hb_util:int(hb_maps:get(<<"to">>, Rewound, -1, Opts))), + ?assertEqual(Canon(101), + hb_maps:get(<<"block-hash">>, Rewound, not_found, Opts)), + %% Global record was actually rewound in the store. + {ok, Persisted} = dev_arweave_scheduler_cache:read_global(Opts), + ?assertEqual(101, + hb_util:int(hb_maps:get(<<"to">>, Persisted, -1, Opts))), + %% The orphan at 102 was overwritten with the canonical block. + {ok, Fixed102} = dev_arweave_scheduler_cache:read_block(102, Opts), + ?assertEqual(Canon(102), + hb_maps:get(<<"indep_hash">>, Fixed102, not_found, Opts)), + %% A divergence deeper than the bound fails closed, untouched. + Deep = #{ <<"from">> => 100, <<"to">> => 102, + <<"block-hash">> => <<"orphan-102">> }, + NeverMatch = fun(H) -> {ok, Mk(H, <<"x-", (hb_util:bin(H))/binary>>, + <<"y">>)} end, + ?assertEqual(not_recovered, + recover_from_reorg(Deep#{ <<"from">> => 102 }, Opts, NeverMatch)), + ok = hb_store:stop(Store). + +reorg_fork_point_test() -> + Canonical = + fun(Height) -> + {ok, <<"canonical-", (hb_util:bin(Height))/binary>>, #{ + <<"height">> => Height + }} + end, + Replace = fun(Height, _Block) -> put({replaced, Height}, true) end, + %% Orphans at 100 and 99; indexed matches canonical at 98. + Indexed = + fun + (98) -> {ok, <<"canonical-98">>}; + (Height) -> {ok, <<"orphan-", (hb_util:bin(Height))/binary>>} + end, + ?assertEqual( + {ok, 98, <<"canonical-98">>}, + fork_point(100, 95, Canonical, Indexed, Replace) + ), + ?assertEqual(true, get({replaced, 100})), + ?assertEqual(true, get({replaced, 99})), + ?assertEqual(undefined, get({replaced, 98})), + %% No agreement within the floor: fail closed. + NeverMatches = fun(_) -> {ok, <<"different">>} end, + ?assertEqual( + not_recovered, + fork_point(100, 97, Canonical, NeverMatches, Replace) + ), + %% A canonical fetch failure aborts recovery rather than guessing. + Failing = fun(_) -> error end, + ?assertEqual( + not_recovered, + fork_point(100, 95, Failing, Indexed, Replace) + ). + block_validation_test() -> State = #{ <<"block-hash">> => <<"previous">> }, Block = From ac486797b116d66b0907a54d7a576c9818d7c8f0 Mon Sep 17 00:00:00 2001 From: xylophonez Date: Fri, 21 Aug 2026 23:23:16 +0100 Subject: [PATCH 2/5] test(arweave-scheduler): add live reorg-injection soak Env-gated (REORG_SOAK=N) soak that seeds an orphaned tail at random fork depths against the live confirmed tip and asserts recover_from_reorg lands exactly on the fork point. Fresh store per iteration. 25/25 correct. --- .../process/dev_arweave_scheduler_sync.erl | 119 ++++++++++++++++++ 1 file changed, 119 insertions(+) diff --git a/src/preloaded/process/dev_arweave_scheduler_sync.erl b/src/preloaded/process/dev_arweave_scheduler_sync.erl index 365acc067..f674a742a 100644 --- a/src/preloaded/process/dev_arweave_scheduler_sync.erl +++ b/src/preloaded/process/dev_arweave_scheduler_sync.erl @@ -1241,6 +1241,125 @@ stop_test_runner(Name) -> hb_name:unregister(Name) end. +reorg_roundtrip_debug_test_() -> + case os:getenv("REORG_DEBUG") of + false -> {"debug disabled", fun() -> ok end}; + _ -> + {timeout, 120, fun() -> + Store = hb_test_utils:test_store( + hb_store_volatile, <<"ar-sched-dbg">>), + ok = hb_store:start(Store), + Opts = #{ <<"store">> => [Store], + <<"scheduler-store">> => [Store], + <<"gateway">> => <<"https://arweave.net">>, + <<"arweave-scheduler-fetch-attempts">> => 3, + <<"priv-wallet">> => ar_wallet:new() }, + {ok, Tip} = confirmed_tip(Opts), + H = Tip - 30, + {ok, B1} = fetch_block_remote(H, Opts), + {ok, B2} = fetch_block_remote(H, Opts), + Hash1 = hb_maps:get(<<"indep_hash">>, B1, none, Opts), + Hash2 = hb_maps:get(<<"indep_hash">>, B2, none, Opts), + ?debugFmt("fetch determinism: ~p", [Hash1 =:= Hash2]), + ?debugFmt("hash1=~p", [Hash1]), + {ok, _} = dev_arweave_scheduler_cache:write_block(H, B1, Opts), + case dev_arweave_scheduler_cache:read_block(H, Opts) of + {ok, B3} -> + Hash3 = + hb_maps:get(<<"indep_hash">>, B3, none, Opts), + ?debugFmt("roundtrip write/read hash match: ~p", + [Hash1 =:= Hash3]), + ?debugFmt("hash3=~p", [Hash3]), + ?debugFmt("B3 keys=~p", + [lists:sort(hb_maps:keys(B3, Opts))]); + Other -> + ?debugFmt("read_block FAILED: ~p", [Other]) + end, + ok = hb_store:stop(Store) + end} + end. + +reorg_soak_test_() -> + case os:getenv("REORG_SOAK") of + false -> {"soak disabled", fun() -> ok end}; + Iters0 -> + Iters = list_to_integer(Iters0), + {timeout, 36000, fun() -> run_reorg_soak(Iters) end} + end. + +run_reorg_soak(Iters) -> + Store0 = hb_test_utils:test_store(hb_store_volatile, <<"ar-sched-soak">>), + ok = hb_store:start(Store0), + Opts0 = #{ <<"store">> => [Store0], <<"scheduler-store">> => [Store0], + <<"gateway">> => <<"https://arweave.net">>, + <<"arweave-scheduler-fetch-attempts">> => 3, + <<"priv-wallet">> => ar_wallet:new() }, + {ok, Tip} = confirmed_tip(Opts0), + ok = hb_store:stop(Store0), + Pass = soak_loop(Iters, 0, Tip, undefined), + ?debugFmt("REORG SOAK: ~p/~p recoveries correct", [Pass, Iters]), + ?assertEqual(Iters, Pass). + +soak_loop(0, Pass, _Tip, _Opts) -> Pass; +soak_loop(N, Pass, Tip, _Ignored) -> + %% Fresh store per iteration so cross-iteration canonical writes cannot + %% shift the fork point. + Store = + hb_test_utils:test_store( + hb_store_volatile, + <<"ar-sched-soak-", (hb_util:bin(N))/binary>>), + ok = hb_store:start(Store), + Opts = #{ <<"store">> => [Store], <<"scheduler-store">> => [Store], + <<"gateway">> => <<"https://arweave.net">>, + <<"arweave-scheduler-fetch-attempts">> => 3, + <<"priv-wallet">> => ar_wallet:new() }, + %% Random fork depth 1..reorg_depth-1 below a live confirmed height. + Depth = 1 + (erlang:phash2({N, erlang:monotonic_time()}, 45)), + Head = Tip - (erlang:phash2(N, 20)), + Fork = Head - Depth, + OK = + case {fetch_block_remote(Fork, Opts), fetch_block_remote(Head, Opts)} of + {{ok, ForkBlk}, {ok, _}} -> + ForkHash = + hb_maps:get(<<"indep_hash">>, ForkBlk, not_found, Opts), + %% Seed: indexed canonical up to Fork, then ORPHAN at Head. + {ok, _} = dev_arweave_scheduler_cache:write_block( + Fork, ForkBlk, Opts), + Orphan = <<"orphan-", (hb_util:bin(N))/binary>>, + {ok, _} = dev_arweave_scheduler_cache:write_block( + Head, + #{ <<"height">> => Head, <<"indep_hash">> => Orphan, + <<"previous_block">> => ForkHash }, + Opts), + State = #{ <<"from">> => 1968888, <<"to">> => Head, + <<"block-hash">> => Orphan }, + {ok, _} = dev_arweave_scheduler_cache:write_global(State, Opts), + case recover_from_reorg(State, Opts) of + {ok, R} -> + case hb_util:int( + hb_maps:get(<<"to">>, R, -1, Opts)) =:= Fork + of + true -> true; + false -> + ?debugFmt("WRONG FORK: got ~p expected ~p", + [hb_maps:get(<<"to">>, R, -1, Opts), Fork]), + false + end; + not_recovered -> skip %% fail-closed (fetch error/floor) + end; + _ -> + %% Network hiccup: don't count against correctness. + skip + end, + ok = hb_store:stop(Store), + case OK of + true -> soak_loop(N - 1, Pass + 1, Tip, undefined); + skip -> timer:sleep(2000), soak_loop(N, Pass, Tip, undefined); + false -> + ?debugFmt("SOAK FAILURE at N=~p head=~p depth=~p", [N, Head, Depth]), + soak_loop(N - 1, Pass, Tip, undefined) + end. + reorg_canonical_fetch_shape_test_() -> {timeout, 60, fun() -> Store = hb_test_utils:test_store( From bdfad3df1415f564e53009ba7a2e195d1bc61733 Mon Sep 17 00:00:00 2001 From: xylophonez Date: Sat, 22 Aug 2026 20:18:18 +0100 Subject: [PATCH 3/5] fix(arweave-scheduler): make reorg recovery parallel, bounded, and non-blocking The reorg-recovery patch placed recover_from_reorg on the synchronous sync path that ~process@1.0/now depends on, and recovered by walking the canonical chain one height at a time via fetch_block_remote (uncached, up to fetch_attempts=20 retries per height) over up to reorg_depth=50 heights. On a flaky network this stalls each recovery attempt for minutes, holds the exclusive sync lock, and starves ~process@1.0/now (zero progress) -- the field symptom. The logic was correct (500/500 soak) but the wall-time was O(depth x attempts x per-request) and a single transient fetch failure aborted the whole recovery. Recovery now fetches the bounded recovery window ([Floor..To]) in ONE concurrent round-trip via fetch_canonical_window/2 (recover_window_workers/2 gives it enough workers to cover the whole window, capped at RECOVERY_MAX_WORKERS), with a small per-height retry budget (recover_fetch_attempts/1 = min(RECOVERY_FETCH_ATTEMPTS=3, fetch_attempts)) so an unreachable gateway fails the window closed fast instead of blocking for minutes. fork_point/4 is now read-only: it locates the deepest matching height from the pre-fetched batch and returns the orphaned heights above it, and only once a valid fork is found does replace_orphans/2 overwrite them and the record rewind -- so the failure path performs no partial write. Correctness guarantees are preserved (deepest-match fork, orphan overwrite above fork, rewind to fork, fail-closed past floor / on unresolvable divergence). Verified: preloaded-store packages cleanly (no device_compile_failed); reorg_fork_point_test passes; on an isolated VPS node normal fresh sync (1.2s), shallow (0.6s) and deep 15-block (1.2s) reorg recovery all succeed and restore orphaned blocks to canonical; a full 51-block window under a black-holed network fails closed in ~15.5s (one 3-attempt round-trip) regardless of depth, versus the old sequential ~103s/height. --- .../process/dev_arweave_scheduler_sync.erl | 257 +++++++++++------- 1 file changed, 165 insertions(+), 92 deletions(-) diff --git a/src/preloaded/process/dev_arweave_scheduler_sync.erl b/src/preloaded/process/dev_arweave_scheduler_sync.erl index f674a742a..eb0983932 100644 --- a/src/preloaded/process/dev_arweave_scheduler_sync.erl +++ b/src/preloaded/process/dev_arweave_scheduler_sync.erl @@ -19,6 +19,8 @@ -define(DEFAULT_FETCH_RETRY_MS, 250). -define(DEFAULT_SYNC_INTERVAL_MS, 1000). -define(DEFAULT_REORG_DEPTH, 50). +-define(RECOVERY_FETCH_ATTEMPTS, 3). +-define(RECOVERY_MAX_WORKERS, 64). -define(PROCESS_IDLE_MS, 1000). %% @doc Bring the global target index to the confirmed chain frontier. The @@ -180,91 +182,155 @@ sync_batch(State = #{ <<"to">> := To }, BatchEnd, Opts) -> %% @doc The indexed tip no longer matches the canonical chain: the node %% indexed a block that the network later orphaned in a reorganization, so -%% every canonical successor fails the previous-hash check forever. Walk the -%% canonical chain downward -- a bounded number of heights, fetched from the -%% network rather than the local cache -- until a canonical block hash matches -%% the indexed block at the same height. Replace each orphaned cached block -%% with its canonical counterpart on the way down so the forward resync -%% cannot re-index the orphan, rewind the global record to the fork point, -%% and let the ordinary forward path re-index the tail. If no indexed block -%% within the bound matches the canonical chain, leave the record untouched -%% and surface the original conflict. +%% every canonical successor fails the previous-hash check forever. Fetch the +%% bounded recovery window (a reorg-depth slice ending at the indexed tip) in +%% one parallel round-trip, always from the network rather than the local +%% cache -- which may still hold the orphans -- then find the deepest height +%% whose canonical block matches the indexed block. Overwrite each orphaned +%% cached block above that fork with its canonical counterpart so the forward +%% resync cannot re-index the orphan, rewind the global record to the fork +%% point, and let the ordinary forward path re-index the tail. If the window +%% cannot be fetched, or no indexed block within the bound matches the +%% canonical chain, leave the record untouched (no partial write) and surface +%% the original conflict. The recovery fetch uses a small retry budget so a +%% stalled or failing gateway fails recovery fast instead of blocking the +%% synchronous sync path for minutes. recover_from_reorg(State, Opts) -> recover_from_reorg( State, Opts, - fun(Height) -> fetch_block_remote(Height, Opts) end + fun(Heights) -> fetch_canonical_window(Heights, Opts) end ). -recover_from_reorg(State = #{ <<"to">> := To, <<"from">> := From }, Opts, FetchBlock) -> +recover_from_reorg(State = #{ <<"to">> := To, <<"from">> := From }, Opts, FetchWindow) -> Floor = max(From, To - reorg_depth(Opts)), - Result = - fork_point( - To, - Floor, - fun(Height) -> - maybe - {ok, Canonical} ?= FetchBlock(Height), - CanonicalHash = - hb_maps:get( - <<"indep_hash">>, Canonical, not_found, Opts - ), - true ?= is_binary(CanonicalHash), - {ok, CanonicalHash, Canonical} - else - _ -> error - end - end, - fun(Height) -> - case dev_arweave_scheduler_cache:read_block(Height, Opts) of - {ok, Indexed} -> - case - hb_maps:get( - <<"indep_hash">>, Indexed, not_found, Opts - ) - of - Hash when is_binary(Hash) -> {ok, Hash}; - _ -> not_found - end; - _ -> not_found - end - end, - fun(Height, Canonical) -> - dev_arweave_scheduler_cache:write_block( - Height, Canonical, Opts + case FetchWindow(lists:seq(Floor, To)) of + {ok, Canonical} -> + case + fork_point( + lists:reverse(Canonical), + Floor, + fun(Height) -> indexed_hash(Height, Opts) end, + Opts ) - end + of + {ok, ForkHeight, ForkHash, Orphans} -> + case replace_orphans(Orphans, Opts) of + ok -> + Rewound = + State#{ + <<"to">> => ForkHeight, + <<"block-hash">> => ForkHash + }, + case + dev_arweave_scheduler_cache:write_global( + Rewound, Opts + ) + of + {ok, _} -> {ok, Rewound}; + _ -> not_recovered + end; + error -> not_recovered + end; + not_recovered -> not_recovered + end; + _ -> not_recovered + end. + +%% @doc Fetch the recovery window concurrently, always from the network (never +%% the local cache, which may still hold the orphans recovery must replace) and +%% with a reduced per-height retry budget so an unreachable gateway fails the +%% whole window closed in one bounded round-trip -- rather than stalling the +%% synchronous sync path, and every caller waiting on it, for minutes. A single +%% unresolved height fails the window closed, leaving the record untouched. +fetch_canonical_window(Heights, Opts) -> + Attempts = recover_fetch_attempts(Opts), + Results = + hb_pmap:parallel_map( + Heights, + fun(Height) -> fetch_block_remote(Height, Attempts, Opts) end, + recover_window_workers(Heights, Opts) ), - case Result of - {ok, ForkHeight, ForkHash} -> - Rewound = - State#{ - <<"to">> => ForkHeight, - <<"block-hash">> => ForkHash - }, - case dev_arweave_scheduler_cache:write_global(Rewound, Opts) of - {ok, _} -> {ok, Rewound}; - _ -> not_recovered + collect_blocks(Heights, Results, []). + +%% @doc Give the bounded recovery window enough workers to fetch in a single +%% concurrent round-trip -- so a stalled gateway fails the whole window in one +%% attempt-budget rather than in `window / block-workers' sequential waves -- +%% while capping the fan-out so a misconfigured reorg depth cannot open an +%% unbounded number of connections. +recover_window_workers(Heights, Opts) -> + Configured = + max( + 1, + hb_util:int( + hb_opts:get( + arweave_scheduler_block_workers, + ?DEFAULT_BLOCK_WORKERS, + Opts + ) + ) + ), + max(Configured, min(length(Heights), ?RECOVERY_MAX_WORKERS)). + +%% @doc The per-height retry budget for recovery's canonical probes, capped +%% well below the forward-sync budget so a flaky network surfaces in seconds +%% rather than stalling recovery for minutes. +recover_fetch_attempts(Opts) -> + max(1, min(?RECOVERY_FETCH_ATTEMPTS, fetch_attempts(Opts))). + +%% @doc Read the locally indexed block's canonical hash at a height, if any. +indexed_hash(Height, Opts) -> + case dev_arweave_scheduler_cache:read_block(Height, Opts) of + {ok, Indexed} -> + case hb_maps:get(<<"indep_hash">>, Indexed, not_found, Opts) of + Hash when is_binary(Hash) -> {ok, Hash}; + _ -> not_found end; - not_recovered -> not_recovered + _ -> not_found end. -%% @doc Find the deepest height, no lower than `Floor', whose canonical -%% network block matches the locally indexed block. Orphaned cached blocks -%% encountered above the fork point are replaced with their canonical -%% counterparts as the walk descends. -fork_point(Height, Floor, _Canonical, _Indexed, _Replace) when Height < Floor -> +%% @doc Find the deepest height, no lower than `Floor', whose canonical block +%% (from the pre-fetched, descending `Batch') matches the locally indexed +%% block. This is read-only: it returns the fork height, its canonical hash, +%% and the orphaned heights above it (each paired with its canonical +%% replacement, ascending) so the caller overwrites orphans only once a valid +%% fork is known -- no writes happen on the failure path. Fails closed when a +%% canonical block lacks a hash or no match exists within the window. +fork_point(Batch, Floor, Indexed, Opts) -> + fork_point(Batch, Floor, Indexed, Opts, []). + +fork_point([], _Floor, _Indexed, _Opts, _Orphans) -> + not_recovered; +fork_point([{Height, _} | _], Floor, _Indexed, _Opts, _Orphans) + when Height < Floor -> not_recovered; -fork_point(Height, Floor, Canonical, Indexed, Replace) -> - case Canonical(Height) of - {ok, CanonicalHash, CanonicalBlock} -> +fork_point([{Height, CanonicalBlock} | Rest], Floor, Indexed, Opts, Orphans) -> + case hb_maps:get(<<"indep_hash">>, CanonicalBlock, not_found, Opts) of + CanonicalHash when is_binary(CanonicalHash) -> case Indexed(Height) of - {ok, CanonicalHash} -> {ok, Height, CanonicalHash}; + {ok, CanonicalHash} -> + {ok, Height, CanonicalHash, Orphans}; _ -> - Replace(Height, CanonicalBlock), - fork_point(Height - 1, Floor, Canonical, Indexed, Replace) + fork_point( + Rest, + Floor, + Indexed, + Opts, + [{Height, CanonicalBlock} | Orphans] + ) end; - error -> not_recovered + _ -> + not_recovered + end. + +%% @doc Overwrite each orphaned cached block above the fork with its canonical +%% counterpart before the record is rewound, so the forward resync cannot +%% re-index an orphan. A failed write fails recovery closed. +replace_orphans([], _Opts) -> ok; +replace_orphans([{Height, Block} | Rest], Opts) -> + case dev_arweave_scheduler_cache:write_block(Height, Block, Opts) of + {ok, _} -> replace_orphans(Rest, Opts); + _ -> error end. %% @doc Return the bounded number of blocks a reorganization recovery may @@ -1409,9 +1475,11 @@ reorg_recovery_integration_test() -> State0 = #{ <<"from">> => 100, <<"to">> => 102, <<"block-hash">> => <<"orphan-102">> }, {ok, _} = dev_arweave_scheduler_cache:write_global(State0, Opts), - %% Canonical fetcher returns the true chain (as the network would). - Fetch = fun(H) -> {ok, Mk(H, Canon(H), - case H of 100 -> <<"canon-99">>; _ -> Canon(H - 1) end)} end, + %% Canonical window fetcher returns the true chain (as the network would) + %% for the whole requested window in one parallel round-trip. + MkBlock = fun(H) -> Mk(H, Canon(H), + case H of 100 -> <<"canon-99">>; _ -> Canon(H - 1) end) end, + Fetch = fun(Heights) -> {ok, [ {H, MkBlock(H)} || H <- Heights ]} end, {ok, Rewound} = recover_from_reorg(State0, Opts, Fetch), %% Fork point is 101 (deepest indexed block matching canonical). ?assertEqual(101, hb_util:int(hb_maps:get(<<"to">>, Rewound, -1, Opts))), @@ -1426,46 +1494,51 @@ reorg_recovery_integration_test() -> ?assertEqual(Canon(102), hb_maps:get(<<"indep_hash">>, Fixed102, not_found, Opts)), %% A divergence deeper than the bound fails closed, untouched. - Deep = #{ <<"from">> => 100, <<"to">> => 102, + Deep = #{ <<"from">> => 102, <<"to">> => 102, <<"block-hash">> => <<"orphan-102">> }, - NeverMatch = fun(H) -> {ok, Mk(H, <<"x-", (hb_util:bin(H))/binary>>, - <<"y">>)} end, + NeverMatch = fun(Heights) -> {ok, [ {H, Mk(H, + <<"x-", (hb_util:bin(H))/binary>>, <<"y">>)} || H <- Heights ]} end, ?assertEqual(not_recovered, - recover_from_reorg(Deep#{ <<"from">> => 102 }, Opts, NeverMatch)), + recover_from_reorg(Deep, Opts, NeverMatch)), + %% A window whose canonical fetch cannot be established fails closed fast, + %% without touching the record. + FailWindow = fun(_Heights) -> + {error, #{ <<"status">> => 503, + <<"reason">> => <<"Arweave block is not retrievable.">> }} + end, + ?assertEqual(not_recovered, + recover_from_reorg(State0, Opts, FailWindow)), ok = hb_store:stop(Store). reorg_fork_point_test() -> - Canonical = - fun(Height) -> - {ok, <<"canonical-", (hb_util:bin(Height))/binary>>, #{ - <<"height">> => Height - }} - end, - Replace = fun(Height, _Block) -> put({replaced, Height}, true) end, + Opts = #{}, + Mk = fun(H) -> #{ <<"height">> => H, + <<"indep_hash">> => <<"canonical-", (hb_util:bin(H))/binary>> } end, + %% Pre-fetched canonical window, descending (highest height first). + Descend = fun(Hi, Lo) -> [ {H, Mk(H)} || H <- lists:seq(Hi, Lo, -1) ] end, %% Orphans at 100 and 99; indexed matches canonical at 98. Indexed = fun (98) -> {ok, <<"canonical-98">>}; (Height) -> {ok, <<"orphan-", (hb_util:bin(Height))/binary>>} end, + %% Deepest match at 98; the orphaned heights above it (99, 100) are + %% reported ascending for replacement, and the fork block itself is not. + %% fork_point is read-only, so it performs no writes. ?assertEqual( - {ok, 98, <<"canonical-98">>}, - fork_point(100, 95, Canonical, Indexed, Replace) + {ok, 98, <<"canonical-98">>, [{99, Mk(99)}, {100, Mk(100)}]}, + fork_point(Descend(100, 95), 95, Indexed, Opts) ), - ?assertEqual(true, get({replaced, 100})), - ?assertEqual(true, get({replaced, 99})), - ?assertEqual(undefined, get({replaced, 98})), %% No agreement within the floor: fail closed. NeverMatches = fun(_) -> {ok, <<"different">>} end, ?assertEqual( not_recovered, - fork_point(100, 97, Canonical, NeverMatches, Replace) + fork_point(Descend(100, 97), 97, NeverMatches, Opts) ), - %% A canonical fetch failure aborts recovery rather than guessing. - Failing = fun(_) -> error end, + %% A canonical block missing its hash fails closed rather than guessing. ?assertEqual( not_recovered, - fork_point(100, 95, Failing, Indexed, Replace) + fork_point([{100, #{ <<"height">> => 100 }}], 100, Indexed, Opts) ). block_validation_test() -> From 65dbfe707045cbf671ab30f732732e0ab4d837ff Mon Sep 17 00:00:00 2001 From: xylophonez Date: Sat, 22 Aug 2026 23:44:57 +0100 Subject: [PATCH 4/5] fix(arweave-scheduler): self-heal stale cached blocks above the frontier A preserved store can hold a stale cached block just above the indexed frontier -- an orphan a previous runtime indexed and left behind -- while the indexed chain itself is canonical. The forward fetch path is cache-first, so every sync pass re-reads the stale block, fails the previous-hash check with the does-not-extend 409, and invokes reorg recovery, which cannot clear it: recovery only repairs cached blocks at or below the frontier, and with the indexed chain already canonical it finds the fork at the frontier itself with nothing to replace. The conflict then recurs forever with zero progress. When the chain check fails, commit_batch now refetches the failing height from the network once, with the bounded recovery retry budget. If the canonical block differs from the batch's copy, the copy was stale: the poisoned cache entry is overwritten and the batch revalidated once. If the canonical block is identical -- a real reorg -- or cannot be fetched, the conflict surfaces unchanged and reorg recovery proceeds exactly as before. One refetch per pass, fail-closed, no partial writes. stale_cache_heal_test reproduces the field scenario offline: canonical chain to the frontier, stale orphan cached above it; the batch heals the cache, validates, and advances without invoking recovery, while identical (real-conflict) and unfetchable refetches surface the original 409 untouched. Verified live on an isolated node: a poisoned frontier+1 heals and resyncs in one pass (0.36s), and a true indexed-orphan reorg still recovers through the unchanged path (0.68s). The preloaded store packages cleanly and the fork-point and recovery integration tests stay green. --- .../process/dev_arweave_scheduler_sync.erl | 142 +++++++++++++++++- 1 file changed, 136 insertions(+), 6 deletions(-) diff --git a/src/preloaded/process/dev_arweave_scheduler_sync.erl b/src/preloaded/process/dev_arweave_scheduler_sync.erl index eb0983932..1fdea2344 100644 --- a/src/preloaded/process/dev_arweave_scheduler_sync.erl +++ b/src/preloaded/process/dev_arweave_scheduler_sync.erl @@ -170,14 +170,76 @@ sync_blocks(State = #{ <<"to">> := To }, Upper, Opts, Recover) -> sync_batch(State = #{ <<"to">> := To }, BatchEnd, Opts) -> Heights = lists:seq(To + 1, BatchEnd), + case fetch_blocks(Heights, Opts) of + {ok, Blocks} -> + commit_batch( + Blocks, + State, + Opts, + fun(Height) -> + fetch_block_remote( + Height, + recover_fetch_attempts(Opts), + Opts + ) + end + ); + Error -> Error + end. + +%% @doc Validate one fetched batch against the indexed frontier and commit +%% it. The forward fetch path is cache-first, so a stale cached block above +%% the frontier -- for example an orphan a previous runtime indexed and left +%% behind in a preserved store -- would otherwise fail the previous-hash +%% check on every pass, forever: reorg recovery cannot clear it, because it +%% only repairs cached blocks at or below the indexed frontier. When the +%% chain check fails, refetch the failing height from the network once, with +%% the bounded recovery retry budget; if the canonical block differs from the +%% batch's copy, the copy was stale: overwrite the poisoned cache entry and +%% revalidate once. If the canonical block is identical to the batch's copy +%% -- or cannot be fetched -- the conflict is real (or undecidable) and +%% surfaces unchanged for reorg recovery to handle. +commit_batch(Blocks, State, Opts, FetchFresh) -> + case validate_blocks(Blocks, State, Opts) of + {ok, NewState} -> + maybe + ok ?= index_blocks(Blocks, Opts), + ok ?= write_blocks(Blocks, Opts), + {ok, _} ?= + dev_arweave_scheduler_cache:write_global(NewState, Opts), + {ok, NewState} + end; + {error, #{ + <<"status">> := 409, + <<"reason">> := + <<"Arweave block does not extend the indexed chain.">>, + <<"block-height">> := Height + }} = Conflict when FetchFresh =/= none -> + case heal_stale_block(Height, Blocks, FetchFresh, Opts) of + {ok, Healed} -> commit_batch(Healed, State, Opts, none); + unchanged -> Conflict + end; + Error -> Error + end. + +%% @doc Replace a batch block that failed the chain check with its freshly +%% fetched canonical counterpart when the two differ, overwriting the +%% poisoned cache entry so the next pass cannot re-read it. Returns +%% `unchanged' when the canonical block matches the batch's copy -- a real +%% conflict, which reorg recovery owns -- or when it cannot be established, +%% failing closed on the original conflict. +heal_stale_block(Height, Blocks, FetchFresh, Opts) -> maybe - {ok, Blocks} ?= fetch_blocks(Heights, Opts), - {ok, NewState} ?= validate_blocks(Blocks, State, Opts), - ok ?= index_blocks(Blocks, Opts), - ok ?= write_blocks(Blocks, Opts), + {Height, Stale} ?= lists:keyfind(Height, 1, Blocks), + {ok, Fresh} ?= FetchFresh(Height), + FreshHash = hb_maps:get(<<"indep_hash">>, Fresh, not_found, Opts), + StaleHash = hb_maps:get(<<"indep_hash">>, Stale, not_found, Opts), + true ?= is_binary(FreshHash) andalso FreshHash =/= StaleHash, {ok, _} ?= - dev_arweave_scheduler_cache:write_global(NewState, Opts), - {ok, NewState} + dev_arweave_scheduler_cache:write_block(Height, Fresh, Opts), + {ok, lists:keyreplace(Height, 1, Blocks, {Height, Fresh})} + else + _ -> unchanged end. %% @doc The indexed tip no longer matches the canonical chain: the node @@ -1455,6 +1517,74 @@ reorg_canonical_fetch_shape_test_() -> ok = hb_store:stop(Store) end}. +%% @doc Field regression: a preserved store holds a STALE cached block just +%% above the frontier (an orphan indexed by an earlier runtime), while the +%% indexed chain itself is canonical. The cache-first forward fetch serves +%% the stale block on every pass and the previous-hash check fails forever; +%% reorg recovery cannot clear it because the divergence is above the +%% frontier and the indexed chain already matches canonical (fork = To, +%% nothing to replace). The batch commit must refetch the failing height, +%% overwrite the poisoned cache entry, revalidate, and advance -- without +%% invoking reorg recovery. +stale_cache_heal_test() -> + Store = hb_test_utils:test_store(hb_store_volatile, <<"ar-sched-heal">>), + ok = hb_store:start(Store), + Opts = #{ <<"store">> => [Store], <<"scheduler-store">> => [Store] }, + Mk = fun(H, Hash, Prev) -> + #{ <<"height">> => H, <<"indep_hash">> => Hash, + <<"previous_block">> => Prev } + end, + Canon = fun(H) -> <<"canon-", (hb_util:bin(H))/binary>> end, + %% Freshly re-indexed canonical chain up to the frontier (To = 101)... + {ok, _} = dev_arweave_scheduler_cache:write_block( + 101, Mk(101, Canon(101), Canon(100)), Opts), + %% ...but the cache still holds a stale orphan ABOVE the frontier whose + %% previous-hash does not extend the indexed chain. + {ok, _} = dev_arweave_scheduler_cache:write_block( + 102, Mk(102, <<"stale-102">>, <<"stale-101">>), Opts), + State0 = #{ <<"from">> => 100, <<"to">> => 101, + <<"block-hash">> => Canon(101) }, + {ok, _} = dev_arweave_scheduler_cache:write_global(State0, Opts), + %% The forward path serves the poison from the cache. + {ok, Blocks} = fetch_blocks([102], Opts), + [{102, Cached}] = Blocks, + ?assertEqual(<<"stale-102">>, + hb_maps:get(<<"indep_hash">>, Cached, not_found, Opts)), + %% The canonical network block (as a fresh refetch would return it). + Fresh = fun(102) -> {ok, Mk(102, Canon(102), Canon(101))} end, + {ok, NewState} = commit_batch(Blocks, State0, Opts, Fresh), + %% The batch healed, validated, and advanced past the poisoned height. + ?assertEqual(102, hb_util:int(hb_maps:get(<<"to">>, NewState, -1, Opts))), + ?assertEqual(Canon(102), + hb_maps:get(<<"block-hash">>, NewState, not_found, Opts)), + %% The poisoned cache entry was overwritten with the canonical block. + {ok, Healed} = dev_arweave_scheduler_cache:read_block(102, Opts), + ?assertEqual(Canon(102), + hb_maps:get(<<"indep_hash">>, Healed, not_found, Opts)), + %% The advance was committed to the global record. + {ok, Persisted} = dev_arweave_scheduler_cache:read_global(Opts), + ?assertEqual(102, + hb_util:int(hb_maps:get(<<"to">>, Persisted, -1, Opts))), + %% A REAL conflict is untouched: when the fresh fetch returns the same + %% block the batch already had, the 409 surfaces for reorg recovery and + %% nothing is written. + {ok, _} = dev_arweave_scheduler_cache:write_block( + 103, Mk(103, <<"orphan-103">>, <<"other-102">>), Opts), + {ok, Blocks2} = fetch_blocks([103], Opts), + Same = fun(103) -> {ok, Mk(103, <<"orphan-103">>, <<"other-102">>)} end, + ?assertMatch( + {error, #{ <<"status">> := 409, <<"block-height">> := 103 }}, + commit_batch(Blocks2, NewState, Opts, Same)), + {ok, Persisted2} = dev_arweave_scheduler_cache:read_global(Opts), + ?assertEqual(102, + hb_util:int(hb_maps:get(<<"to">>, Persisted2, -1, Opts))), + %% An unfetchable fresh block also fails closed on the original conflict. + Failing = fun(_) -> {error, #{ <<"status">> => 503 }} end, + ?assertMatch( + {error, #{ <<"status">> := 409 }}, + commit_batch(Blocks2, NewState, Opts, Failing)), + ok = hb_store:stop(Store). + reorg_recovery_integration_test() -> Store = hb_test_utils:test_store(hb_store_volatile, <<"ar-sched-reorg">>), ok = hb_store:start(Store), From d08cb1b835eebfade347619d68bf163e0e81c3ec Mon Sep 17 00:00:00 2001 From: xylophonez Date: Sun, 23 Aug 2026 00:06:30 +0100 Subject: [PATCH 5/5] fix(arweave-scheduler): heal the whole batch when the cache is stale The stale-cache heal keyed on the conflict height alone cannot clear an orphan RUN: a preserved store can cache an orphan that EXTENDS the canonical frontier -- passing the previous-hash check -- with the divergence only surfacing one height later. Refetching just the conflict height then yields the canonical block, which still fails to extend the stale orphan ahead of it in the batch, and the pass wedges on the 409 forever exactly as before. On a does-not-extend conflict, commit_batch now distrusts the cache for the whole batch: every height is refetched remote-only in one bounded parallel round-trip (the recovery window fetcher; a batch is at most the block-batch size), each cached entry whose canonical hash differs is overwritten, the batch is rebuilt from the fresh blocks, and validation runs once more. If the fully-remote batch still conflicts, the divergence is the indexed chain's own -- a genuine reorganization -- and the conflict surfaces unchanged for reorg recovery. A refetch that cannot be established, misaligns, or fails to write fails closed on the original conflict, and the global record is never touched on a failure path. stale_cache_run_heal_test reproduces the field shape offline: a canonical indexed chain with a two-block stale orphan run above the frontier whose first orphan extends the canonical tip; the batch heal replaces both poisoned entries and advances without invoking recovery. The single-block heal, real-conflict pass-through, and fetch-failure fail-closed cases stay green, as do the fork-point and recovery integration tests. Verified live on an isolated node: the injected orphan run heals and resyncs in one pass (0.60s, both cache entries restored to canonical), and a true indexed-orphan reorg still recovers through the unchanged path (0.78s). The preloaded store packages cleanly. --- .../process/dev_arweave_scheduler_sync.erl | 172 +++++++++++++----- 1 file changed, 127 insertions(+), 45 deletions(-) diff --git a/src/preloaded/process/dev_arweave_scheduler_sync.erl b/src/preloaded/process/dev_arweave_scheduler_sync.erl index 1fdea2344..3b5fa259e 100644 --- a/src/preloaded/process/dev_arweave_scheduler_sync.erl +++ b/src/preloaded/process/dev_arweave_scheduler_sync.erl @@ -176,29 +176,27 @@ sync_batch(State = #{ <<"to">> := To }, BatchEnd, Opts) -> Blocks, State, Opts, - fun(Height) -> - fetch_block_remote( - Height, - recover_fetch_attempts(Opts), - Opts - ) - end + fun(Hs) -> fetch_canonical_window(Hs, Opts) end ); Error -> Error end. %% @doc Validate one fetched batch against the indexed frontier and commit -%% it. The forward fetch path is cache-first, so a stale cached block above -%% the frontier -- for example an orphan a previous runtime indexed and left -%% behind in a preserved store -- would otherwise fail the previous-hash -%% check on every pass, forever: reorg recovery cannot clear it, because it -%% only repairs cached blocks at or below the indexed frontier. When the -%% chain check fails, refetch the failing height from the network once, with -%% the bounded recovery retry budget; if the canonical block differs from the -%% batch's copy, the copy was stale: overwrite the poisoned cache entry and -%% revalidate once. If the canonical block is identical to the batch's copy -%% -- or cannot be fetched -- the conflict is real (or undecidable) and -%% surfaces unchanged for reorg recovery to handle. +%% it. The forward fetch path is cache-first, so stale cached blocks above +%% the frontier -- for example an orphan run a previous runtime indexed and +%% left behind in a preserved store -- would otherwise fail the previous-hash +%% check on every pass, forever: reorg recovery cannot clear them, because it +%% only repairs cached blocks at or below the indexed frontier. A stale run +%% can even begin with an orphan that extends the canonical frontier, so the +%% conflict height alone does not locate the poison. When the chain check +%% fails, distrust the cache for the whole batch: refetch every height from +%% the network in one bounded parallel round-trip, overwrite each cached +%% entry whose canonical hash differs, rebuild the batch from the fresh +%% blocks, and revalidate once. If the fully-remote batch still conflicts, +%% the indexed chain itself has diverged -- a genuine reorganization -- and +%% the conflict surfaces unchanged for reorg recovery to handle. If the +%% refetch cannot be established, fail closed on the original conflict; the +%% global record is never touched on a failure path. commit_batch(Blocks, State, Opts, FetchFresh) -> case validate_blocks(Blocks, State, Opts) of {ok, NewState} -> @@ -213,35 +211,57 @@ commit_batch(Blocks, State, Opts, FetchFresh) -> <<"status">> := 409, <<"reason">> := <<"Arweave block does not extend the indexed chain.">>, - <<"block-height">> := Height + <<"block-height">> := _Height }} = Conflict when FetchFresh =/= none -> - case heal_stale_block(Height, Blocks, FetchFresh, Opts) of + case heal_stale_batch(Blocks, FetchFresh, Opts) of {ok, Healed} -> commit_batch(Healed, State, Opts, none); unchanged -> Conflict end; Error -> Error end. -%% @doc Replace a batch block that failed the chain check with its freshly -%% fetched canonical counterpart when the two differ, overwriting the -%% poisoned cache entry so the next pass cannot re-read it. Returns -%% `unchanged' when the canonical block matches the batch's copy -- a real -%% conflict, which reorg recovery owns -- or when it cannot be established, -%% failing closed on the original conflict. -heal_stale_block(Height, Blocks, FetchFresh, Opts) -> - maybe - {Height, Stale} ?= lists:keyfind(Height, 1, Blocks), - {ok, Fresh} ?= FetchFresh(Height), - FreshHash = hb_maps:get(<<"indep_hash">>, Fresh, not_found, Opts), - StaleHash = hb_maps:get(<<"indep_hash">>, Stale, not_found, Opts), - true ?= is_binary(FreshHash) andalso FreshHash =/= StaleHash, - {ok, _} ?= - dev_arweave_scheduler_cache:write_block(Height, Fresh, Opts), - {ok, lists:keyreplace(Height, 1, Blocks, {Height, Fresh})} - else +%% @doc Refetch the conflicted batch's heights remote-only and replace the +%% batch with the canonical blocks, overwriting each poisoned cache entry +%% whose hash differs so the next pass cannot re-read it. Returns `unchanged' +%% when every canonical block already matches the batch's copy -- the batch +%% was not the problem, so the conflict is the indexed chain's and reorg +%% recovery owns it -- or when the canonical batch cannot be established or +%% written, failing closed on the original conflict. +heal_stale_batch(Blocks, FetchFresh, Opts) -> + case FetchFresh([ Height || {Height, _} <- Blocks ]) of + {ok, Fresh} -> + case stale_pairs(Blocks, Fresh, Opts) of + [] -> unchanged; + mismatch -> unchanged; + Stale -> + case replace_orphans(Stale, Opts) of + ok -> {ok, Fresh}; + error -> unchanged + end + end; _ -> unchanged end. +%% @doc Pair the batch with its freshly fetched canonical counterparts and +%% keep the heights whose canonical block differs from the batch's copy. +%% Returns `mismatch' -- aborting the heal -- when the canonical batch does +%% not align height-for-height or a canonical block lacks its hash. +stale_pairs([], [], _Opts) -> []; +stale_pairs([{Height, Cached} | Blocks], [{Height, Fresh} | Rest], Opts) -> + FreshHash = hb_maps:get(<<"indep_hash">>, Fresh, not_found, Opts), + CachedHash = hb_maps:get(<<"indep_hash">>, Cached, not_found, Opts), + case is_binary(FreshHash) of + false -> mismatch; + true -> + case stale_pairs(Blocks, Rest, Opts) of + mismatch -> mismatch; + Tail when FreshHash =/= CachedHash -> + [{Height, Fresh} | Tail]; + Tail -> Tail + end + end; +stale_pairs(_Blocks, _Fresh, _Opts) -> mismatch. + %% @doc The indexed tip no longer matches the canonical chain: the node %% indexed a block that the network later orphaned in a reorganization, so %% every canonical successor fails the previous-hash check forever. Fetch the @@ -385,9 +405,11 @@ fork_point([{Height, CanonicalBlock} | Rest], Floor, Indexed, Opts, Orphans) -> not_recovered end. -%% @doc Overwrite each orphaned cached block above the fork with its canonical -%% counterpart before the record is rewound, so the forward resync cannot -%% re-index an orphan. A failed write fails recovery closed. +%% @doc Overwrite each orphaned cached block with its canonical counterpart, +%% so no later pass can re-read stale state: reorg recovery clears the blocks +%% above the fork before the record is rewound, and the batch heal clears +%% poisoned entries above the frontier before revalidation. A failed write +%% fails the caller closed. replace_orphans([], _Opts) -> ok; replace_orphans([{Height, Block} | Rest], Opts) -> case dev_arweave_scheduler_cache:write_block(Height, Block, Opts) of @@ -1550,8 +1572,8 @@ stale_cache_heal_test() -> [{102, Cached}] = Blocks, ?assertEqual(<<"stale-102">>, hb_maps:get(<<"indep_hash">>, Cached, not_found, Opts)), - %% The canonical network block (as a fresh refetch would return it). - Fresh = fun(102) -> {ok, Mk(102, Canon(102), Canon(101))} end, + %% The canonical network batch (as a fresh refetch would return it). + Fresh = fun([102]) -> {ok, [{102, Mk(102, Canon(102), Canon(101))}]} end, {ok, NewState} = commit_batch(Blocks, State0, Opts, Fresh), %% The batch healed, validated, and advanced past the poisoned height. ?assertEqual(102, hb_util:int(hb_maps:get(<<"to">>, NewState, -1, Opts))), @@ -1566,25 +1588,85 @@ stale_cache_heal_test() -> ?assertEqual(102, hb_util:int(hb_maps:get(<<"to">>, Persisted, -1, Opts))), %% A REAL conflict is untouched: when the fresh fetch returns the same - %% block the batch already had, the 409 surfaces for reorg recovery and + %% batch the cache already held, the 409 surfaces for reorg recovery and %% nothing is written. {ok, _} = dev_arweave_scheduler_cache:write_block( 103, Mk(103, <<"orphan-103">>, <<"other-102">>), Opts), {ok, Blocks2} = fetch_blocks([103], Opts), - Same = fun(103) -> {ok, Mk(103, <<"orphan-103">>, <<"other-102">>)} end, + Same = fun([103]) -> {ok, [{103, Mk(103, <<"orphan-103">>, <<"other-102">>)}]} end, ?assertMatch( {error, #{ <<"status">> := 409, <<"block-height">> := 103 }}, commit_batch(Blocks2, NewState, Opts, Same)), {ok, Persisted2} = dev_arweave_scheduler_cache:read_global(Opts), ?assertEqual(102, hb_util:int(hb_maps:get(<<"to">>, Persisted2, -1, Opts))), - %% An unfetchable fresh block also fails closed on the original conflict. + %% An unfetchable fresh batch also fails closed on the original conflict. Failing = fun(_) -> {error, #{ <<"status">> => 503 }} end, ?assertMatch( {error, #{ <<"status">> := 409 }}, commit_batch(Blocks2, NewState, Opts, Failing)), ok = hb_store:stop(Store). +%% @doc Field regression: the stale cache holds an orphan RUN, and its first +%% orphan EXTENDS the canonical frontier -- so the previous-hash check passes +%% at the first poisoned height and the conflict only surfaces one height +%% later, where the run diverges. A heal keyed on the conflict height alone +%% refetches the wrong block and re-conflicts forever; the batch heal must +%% distrust the whole batch, replace every stale entry, and advance without +%% invoking reorg recovery. +stale_cache_run_heal_test() -> + Store = hb_test_utils:test_store( + hb_store_volatile, <<"ar-sched-heal-run">>), + ok = hb_store:start(Store), + Opts = #{ <<"store">> => [Store], <<"scheduler-store">> => [Store] }, + Mk = fun(H, Hash, Prev) -> + #{ <<"height">> => H, <<"indep_hash">> => Hash, + <<"previous_block">> => Prev } + end, + Canon = fun(H) -> <<"canon-", (hb_util:bin(H))/binary>> end, + %% Canonical indexed chain to the frontier (To = 201)... + {ok, _} = dev_arweave_scheduler_cache:write_block( + 201, Mk(201, Canon(201), Canon(200)), Opts), + %% ...and a stale 2-block orphan run above it: the orphan at 202 extends + %% the CANONICAL frontier, so it validates; the run diverges at 203. + {ok, _} = dev_arweave_scheduler_cache:write_block( + 202, Mk(202, <<"orphan-202">>, Canon(201)), Opts), + {ok, _} = dev_arweave_scheduler_cache:write_block( + 203, Mk(203, <<"orphan-203">>, Canon(202)), Opts), + State0 = #{ <<"from">> => 200, <<"to">> => 201, + <<"block-hash">> => Canon(201) }, + {ok, _} = dev_arweave_scheduler_cache:write_global(State0, Opts), + %% The forward path serves the poisoned run from the cache. + {ok, Blocks} = fetch_blocks([202, 203], Opts), + [{202, Cached202}, {203, _}] = Blocks, + ?assertEqual(<<"orphan-202">>, + hb_maps:get(<<"indep_hash">>, Cached202, not_found, Opts)), + %% The canonical network batch. + Fresh = + fun([202, 203]) -> + {ok, [ + {202, Mk(202, Canon(202), Canon(201))}, + {203, Mk(203, Canon(203), Canon(202))} + ]} + end, + {ok, NewState} = commit_batch(Blocks, State0, Opts, Fresh), + %% The whole run healed and the batch advanced past it. + ?assertEqual(203, hb_util:int(hb_maps:get(<<"to">>, NewState, -1, Opts))), + ?assertEqual(Canon(203), + hb_maps:get(<<"block-hash">>, NewState, not_found, Opts)), + %% BOTH poisoned cache entries were overwritten with canonical blocks. + {ok, Healed202} = dev_arweave_scheduler_cache:read_block(202, Opts), + ?assertEqual(Canon(202), + hb_maps:get(<<"indep_hash">>, Healed202, not_found, Opts)), + {ok, Healed203} = dev_arweave_scheduler_cache:read_block(203, Opts), + ?assertEqual(Canon(203), + hb_maps:get(<<"indep_hash">>, Healed203, not_found, Opts)), + %% The advance was committed to the global record. + {ok, Persisted} = dev_arweave_scheduler_cache:read_global(Opts), + ?assertEqual(203, + hb_util:int(hb_maps:get(<<"to">>, Persisted, -1, Opts))), + ok = hb_store:stop(Store). + reorg_recovery_integration_test() -> Store = hb_test_utils:test_store(hb_store_volatile, <<"ar-sched-reorg">>), ok = hb_store:start(Store),