diff --git a/src/preloaded/process/dev_arweave_scheduler_sync.erl b/src/preloaded/process/dev_arweave_scheduler_sync.erl index b11719d4f..3b5fa259e 100644 --- a/src/preloaded/process/dev_arweave_scheduler_sync.erl +++ b/src/preloaded/process/dev_arweave_scheduler_sync.erl @@ -18,6 +18,9 @@ -define(DEFAULT_FETCH_ATTEMPTS, 20). -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 @@ -129,9 +132,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,17 +153,285 @@ 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), - {ok, NewState} ?= validate_blocks(Blocks, State, Opts), - ok ?= index_blocks(Blocks, Opts), - ok ?= write_blocks(Blocks, Opts), - {ok, _} ?= - dev_arweave_scheduler_cache:write_global(NewState, Opts), - sync_blocks(NewState, Upper, Opts) + case fetch_blocks(Heights, Opts) of + {ok, Blocks} -> + commit_batch( + Blocks, + State, + Opts, + 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 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} -> + 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_batch(Blocks, FetchFresh, Opts) of + {ok, Healed} -> commit_batch(Healed, State, Opts, none); + unchanged -> Conflict + end; + Error -> Error end. +%% @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 +%% 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(Heights) -> fetch_canonical_window(Heights, Opts) end + ). + +recover_from_reorg(State = #{ <<"to">> := To, <<"from">> := From }, Opts, FetchWindow) -> + Floor = max(From, To - reorg_depth(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 + ) + 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) + ), + 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_found + end. + +%% @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, 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, Orphans}; + _ -> + fork_point( + Rest, + Floor, + Indexed, + Opts, + [{Height, CanonicalBlock} | Orphans] + ) + end; + _ -> + not_recovered + end. + +%% @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 + {ok, _} -> replace_orphans(Rest, Opts); + _ -> error + 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 +1391,368 @@ 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( + 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}. + +%% @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 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))), + ?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 + %% 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, [{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 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), + 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 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))), + ?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">> => 102, <<"to">> => 102, + <<"block-hash">> => <<"orphan-102">> }, + NeverMatch = fun(Heights) -> {ok, [ {H, Mk(H, + <<"x-", (hb_util:bin(H))/binary>>, <<"y">>)} || H <- Heights ]} end, + ?assertEqual(not_recovered, + 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() -> + 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">>, [{99, Mk(99)}, {100, Mk(100)}]}, + fork_point(Descend(100, 95), 95, Indexed, Opts) + ), + %% No agreement within the floor: fail closed. + NeverMatches = fun(_) -> {ok, <<"different">>} end, + ?assertEqual( + not_recovered, + fork_point(Descend(100, 97), 97, NeverMatches, Opts) + ), + %% A canonical block missing its hash fails closed rather than guessing. + ?assertEqual( + not_recovered, + fork_point([{100, #{ <<"height">> => 100 }}], 100, Indexed, Opts) + ). + block_validation_test() -> State = #{ <<"block-hash">> => <<"previous">> }, Block =