From 8cdc8d20a5e05e78da6165348cd3eef84600b0f0 Mon Sep 17 00:00:00 2001 From: Sam Williams Date: Thu, 16 Jul 2026 14:14:00 -0400 Subject: [PATCH 1/7] refactor(copycat): share the local item-caching helper Caching an L1 transaction header (every mode) and caching a fully-parsed bundle item (full mode) used the same scope-convert-write sequence with different codecs. Share it as `cache_item/3'; behaviour is unchanged (the header path keeps its log-and-skip guard, the full-mode path still asserts the write). --- src/preloaded/query/dev_copycat_arweave.erl | 19 ++++++++++--------- 1 file changed, 10 insertions(+), 9 deletions(-) diff --git a/src/preloaded/query/dev_copycat_arweave.erl b/src/preloaded/query/dev_copycat_arweave.erl index 0b06b0666..5b5b6675d 100644 --- a/src/preloaded/query/dev_copycat_arweave.erl +++ b/src/preloaded/query/dev_copycat_arweave.erl @@ -736,10 +736,7 @@ index_full_bundle_items( ok = case {IndexMode, ParseResult} of {full, {ok, _, Parsed}} -> - LocalOpts = hb_store:scope(Opts, local), - Msg = hb_message:convert( - Parsed, <<"structured@1.0">>, <<"ans104@1.0">>, LocalOpts), - {ok, _Path} = hb_cache:write(Msg, LocalOpts), + {ok, _Path} = cache_item(Parsed, <<"ans104@1.0">>, Opts), ok; _ -> ok end, @@ -876,12 +873,8 @@ cache_tx_header(TX, Opts) -> end. write_tx_header(TX, Opts) -> - LocalOpts = hb_store:scope(Opts, local), try - Msg = - hb_message:convert( - TX, <<"structured@1.0">>, <<"tx@1.0">>, LocalOpts), - {ok, _} = hb_cache:write(Msg, LocalOpts), + {ok, _} = cache_item(TX, <<"tx@1.0">>, Opts), ok catch Class:Reason -> @@ -896,6 +889,14 @@ write_tx_header(TX, Opts) -> ok end. +%% @doc Cache an item decoded during the block scan in the local store as a +%% structured message, keyed by the codec it is committed with. Its fields +%% (notably `target') then become locally matchable. +cache_item(Item, Codec, Opts) -> + LocalOpts = hb_store:scope(Opts, local), + Msg = hb_message:convert(Item, <<"structured@1.0">>, Codec, LocalOpts), + hb_cache:write(Msg, LocalOpts). + %% @doc Record event metrics (count and duration) using hb_event:record. record_event_metrics(MetricName, Count, Duration) -> hb_event:record(<<"arweave_block_count">>, MetricName, #{}, Count), From e14ba678e2c6e06f14ea37e6829b213bd90f4114 Mon Sep 17 00:00:00 2001 From: Sam Williams Date: Thu, 16 Jul 2026 14:11:27 -0400 Subject: [PATCH 2/7] refactor(scheduler): extract the logic shared by scheduler devices into lib_scheduler `~scheduler@1.0' and the `~arweave-scheduler@1.0' device added later in this branch present the same consumer interface, so everything about interpreting a scheduling request -- and storing and serving its results -- is identical regardless of how the schedule is sourced. Move that shared logic into a new `lib_scheduler' device library: - `find_target_id/4,3' and `find_message_to_schedule/3' (request parsing), `load_message_to_schedule/3' (find + fully load, 404 on failure), and `only_committed/2' (reduce to committed components, 400 on failure). - `parse_slot_range/2' (the `from'/`to' slot parse; `from' is now clamped non-negative after coercion, which also catches a negative binary input) and `read_assignment_range/5' (bounded range read + truncation flag), with the single `?MAX_ASSIGNMENT_QUERY_LEN' definition and `max_assignment_query_len/0' for callers that bound remote requests by it. - The assignment cache: `cache_opts/1', `write_assignment/3' and `read_assignment/4', parameterized by each device's pseudo-path prefix. The read restores unsigned commitment IDs and converts legacy `ao.TN.1' assignments, exactly as `dev_scheduler_cache:read/3' did. - `base_assignment/4', the assignment keys every scheduler shares -- `dev_scheduler_server:do_assign' merges its block fields and commits; other schedulers merge their own position fields. - `at_slot/2', `slot_unavailable/0' and `format_opts/1' (the no-cache/no-await formatting options, previously inlined in several places). `dev_scheduler', `dev_scheduler_cache', `dev_scheduler_formats' and `dev_scheduler_server' delegate to the library; no behavioural change. --- src/preloaded/process/dev_scheduler.erl | 240 +++---------- src/preloaded/process/dev_scheduler_cache.erl | 96 +----- .../process/dev_scheduler_formats.erl | 7 +- .../process/dev_scheduler_server.erl | 22 +- src/preloaded/process/lib_scheduler.erl | 322 ++++++++++++++++++ 5 files changed, 384 insertions(+), 303 deletions(-) create mode 100644 src/preloaded/process/lib_scheduler.erl diff --git a/src/preloaded/process/dev_scheduler.erl b/src/preloaded/process/dev_scheduler.erl index fe8dff675..0333080b9 100644 --- a/src/preloaded/process/dev_scheduler.erl +++ b/src/preloaded/process/dev_scheduler.erl @@ -15,7 +15,7 @@ %%% -module(dev_scheduler). --device_libraries([lib_process]). +-device_libraries([lib_process, lib_scheduler]). %%% AO-Core API functions: -export([info/0]). %%% Local scheduling functions: @@ -30,8 +30,6 @@ -include("include/hb.hrl"). -include_lib("eunit/include/eunit.hrl"). -%%% The maximum number of assignments that we will query/return at a time. --define(MAX_ASSIGNMENT_QUERY_LEN, 1000). %%% The timeout for a lookahead worker. -define(LOOKAHEAD_TIMEOUT, 1500). @@ -84,14 +82,7 @@ next(Base, Req, Opts) -> ?event(next, started_next), ?event(next_profiling, started_next), Schedule = message_cached_assignments(Base, Opts), - LastProcessed = - hb_util:int( - hb_ao:get( - <<"at-slot">>, - Base, - Opts#{ <<"hashpath">> => ignore } - ) - ), + LastProcessed = lib_scheduler:at_slot(Base, Opts), ?event(next_profiling, got_last_processed), ?event(debug_next, {in_message_cache, {schedule, Schedule}}), ?event(next, {last_processed, LastProcessed, {message_cache, length(Schedule)}}), @@ -111,12 +102,7 @@ next(Base, Req, Opts) -> ?event(next_profiling, got_no_assignments), {error, Reason}; {ok, [], _} -> - {error, #{ - <<"status">> => 404, - <<"reason">> => - <<"Requested slot not yet available in schedule.">> - } - }; + lib_scheduler:slot_unavailable(); {ok, Assignments, Lookahead} -> ?event(next_profiling, got_assignments), validate_next_slot(Base, Assignments, Lookahead, LastProcessed, Opts) @@ -375,21 +361,9 @@ schedule(Base, Req, Opts) -> %% for this scheduler. If so, it schedules the message and returns the assignment. post_schedule(Base, Req, Opts) -> ?event(scheduling_message), - % Find the target message to schedule: - RawToSched = find_message_to_schedule(Base, Req, Opts), - % If the message can not be properly loaded, this will throw an error - % before scheduling the message. - try hb_cache:ensure_all_loaded(RawToSched, Opts) of - ToSched -> - do_post_schedule(Base, Req, ToSched, Opts) - catch - error:{necessary_message_not_found, _, _} -> - {error, - #{ - <<"status">> => 404, - <<"body">> => <<"Cannot fully load message to schedule.">> - } - } + maybe + {ok, ToSched} ?= lib_scheduler:load_message_to_schedule(Base, Req, Opts), + do_post_schedule(Base, Req, ToSched, Opts) end. do_post_schedule(Base, Req, ToSched, Opts) -> @@ -397,53 +371,44 @@ do_post_schedule(Base, Req, ToSched, Opts) -> % Find the ProcessID of the target message: % - If it is a Process, use the ID of the message. % - If not, use the target as the ProcessID. - ProcID = find_target_id(Base, Req, ToSched, Opts), + ProcID = lib_scheduler:find_target_id(Base, Req, ToSched, Opts), ?event({proc_id, ProcID}), % Filter all unsigned keys from the source message. - case hb_message:with_only_committed(ToSched, Opts) of - {ok, OnlyCommitted} -> - ?event( - {post_schedule, - {schedule_id, ProcID}, - {message, ToSched} - } - ), - % Find the relevant scheduler server for the given process and - % message, start a new one if necessary, or return a redirect to the - % correct remote scheduler. - case find_server(ProcID, Base, ToSched, Opts) of - {local, PID} -> - ?event({scheduling_locally, {proc_id, ProcID}, {pid, PID}}), - post_local_schedule(ProcID, PID, OnlyCommitted, Opts); - {redirect, Redirect} -> - ?event({process_is_remote, {redirect, Redirect}}), - case hb_opts:get(scheduler_follow_redirects, true, Opts) of - true -> - ?event({proxying_to_remote_scheduler, - {redirect, Redirect}, - {msg, OnlyCommitted} - }), - post_remote_schedule( - ProcID, - Redirect, - OnlyCommitted, - Opts - ); - false -> {ok, Redirect} - end; - {error, Error} -> - ?event({error_finding_scheduler, {error, Error}}), - {error, Error} - end; - {error, Err} -> - {error, - #{ - <<"status">> => 400, - <<"body">> => <<"Message invalid: ", - "Committed components cannot be validated.">>, - <<"reason">> => Err - } + maybe + {ok, OnlyCommitted} ?= lib_scheduler:only_committed(ToSched, Opts), + ?event( + {post_schedule, + {schedule_id, ProcID}, + {message, ToSched} } + ), + % Find the relevant scheduler server for the given process and + % message, start a new one if necessary, or return a redirect to the + % correct remote scheduler. + case find_server(ProcID, Base, ToSched, Opts) of + {local, PID} -> + ?event({scheduling_locally, {proc_id, ProcID}, {pid, PID}}), + post_local_schedule(ProcID, PID, OnlyCommitted, Opts); + {redirect, Redirect} -> + ?event({process_is_remote, {redirect, Redirect}}), + case hb_opts:get(scheduler_follow_redirects, true, Opts) of + true -> + ?event({proxying_to_remote_scheduler, + {redirect, Redirect}, + {msg, OnlyCommitted} + }), + post_remote_schedule( + ProcID, + Redirect, + OnlyCommitted, + Opts + ); + false -> {ok, Redirect} + end; + {error, Error} -> + ?event({error_finding_scheduler, {error, Error}}), + {error, Error} + end end. %% @doc Post schedule the message. `Req' by this point has been refined to only @@ -717,7 +682,7 @@ find_remote_scheduler(ProcID, Scheduler, Opts) -> %% @doc Returns information about the current slot for a process. slot(M1, M2, Opts) -> ?event({getting_current_slot, {msg, M1}}), - ProcID = find_target_id(M1, M2, Opts), + ProcID = lib_scheduler:find_target_id(M1, M2, Opts), case find_server(ProcID, M1, Opts) of {local, PID} -> ?event({getting_current_slot, {proc_id, ProcID}}), @@ -813,18 +778,8 @@ remote_slot(<<"ao.TN.1">>, ProcID, Node, Opts) -> %% two slots -- labelled as `from' and `to'. If the schedule is not local, %% we redirect to the remote scheduler or proxy based on the node opts. get_schedule(Base, Req, Opts) -> - ProcID = hb_util:human_id(find_target_id(Base, Req, Opts)), - From = - case hb_ao:get(<<"from">>, Req, not_found, Opts) of - not_found -> 0; - X when X < 0 -> 0; - FromRes -> hb_util:int(FromRes) - end, - To = - case hb_ao:get(<<"to">>, Req, not_found, Opts) of - not_found -> undefined; - ToRes -> hb_util:int(ToRes) - end, + ProcID = hb_util:human_id(lib_scheduler:find_target_id(Base, Req, Opts)), + {From, To} = lib_scheduler:parse_slot_range(Req, Opts), Format = hb_ao:get(<<"accept">>, Req, <<"application/http">>, Opts), ?event( {parsed_get_schedule, @@ -963,7 +918,7 @@ do_get_remote_schedule(ProcID, LocalAssignments, From, To, Redirect, Opts) -> << ProcID/binary, "?process-id=", ProcID/binary, FromBin/binary, ToParam/binary, - "&limit=", (hb_util:bin(?MAX_ASSIGNMENT_QUERY_LEN))/binary + "&limit=", (hb_util:bin(lib_scheduler:max_assignment_query_len()))/binary >> end, ?event({getting_remote_schedule, {node, {string, Node}}, {path, {string, Path}}}), @@ -1295,80 +1250,6 @@ post_legacy_schedule(ProcID, OnlyCommitted, Node, Opts) -> %%% Private methods -%% @doc Find the schedule ID from a given request. The precidence order for -%% search is as follows: -%% 1. `ToSched/id' when `ToSched' has `type: Process' -%% 2. `ToSched/target' when `ToSched' has a `target' key -%% 2. `Req/target' -%% 3. `Req/id' when `Req' has `type: Process' -%% 4. `Base/process/id' -%% 5. `Base/id' when `Base' has `type: Process' -%% 6. `Req/id' -find_target_id(Base, Req, ToSched, Opts) -> - case hb_ao:get(<<"type">>, ToSched, not_found, Opts) of - <<"Process">> -> - lib_process:process_id(ToSched, #{}, Opts); - _ -> - case hb_ao:get(<<"target">>, ToSched, not_found, Opts) of - not_found -> find_target_id(Base, Req, Opts); - Target -> hb_util:human_id(Target) - end - end. -find_target_id(Base, Req, Opts) -> - TempOpts = Opts#{ <<"hashpath">> => ignore }, - Res = case hb_ao:resolve(Req, <<"target">>, TempOpts) of - {ok, Target} -> - % ID found at Req/target - Target; - _ -> - case hb_ao:resolve(Req, <<"type">>, TempOpts) of - {ok, <<"Process">>} -> - % Req is a Process, so the ID is at Req/id - lib_process:process_id(Req, #{}, Opts); - _ -> - case hb_ao:resolve(Base, <<"process">>, TempOpts) of - {ok, _Process} -> - lib_process:process_id(Base, #{}, Opts); - _ -> - % Does the message have a type of process? - case hb_ao:get(<<"type">>, Base, TempOpts) of - <<"Process">> -> - % Yes: Base is the process. - lib_process:process_id(Base, #{}, Opts); - _ -> - % No: Req is the target process. - lib_process:process_id(Req, #{}, Opts) - end - end - end - end, - ?event({found_id, {id, Res}, {base, Base}, {req, Req}}), - Res. - -%% @doc Search the given base and request message pair to find the message to -%% schedule. The precidence order for search is as follows: -%% 1. A key in `Req' with the value `self', indicating that the entire message -%% is the subject. -%% 2. A key in `Req' with another value, present in that message. -%% 3. The body of the message. -%% 4. The message itself. -find_message_to_schedule(Base, Req, Opts) -> - Subject = - hb_ao:get( - <<"subject">>, - Req, - not_found, - Opts#{ <<"hashpath">> => ignore } - ), - case Subject of - <<"base">> -> Base; - <<"self">> -> Req; - not_found -> - hb_ao:get(<<"body">>, Req, Req, Opts#{ <<"hashpath">> => ignore }); - Subject -> - hb_ao:get(Subject, Req, Opts#{ <<"hashpath">> => ignore }) - end. - %% @doc Generate a `GET /schedule' response for a process. generate_local_schedule(Format, ProcID, From, To, Opts) -> ?event( @@ -1405,35 +1286,8 @@ get_local_assignments(ProcID, From, undefined, Opts) -> end; get_local_assignments(ProcID, From, RequestedTo, Opts) -> ?event({handling_req_to_get_assignments, ProcID, {from, From}, {to, RequestedTo}}), - ComputedTo = - case (RequestedTo - From) > ?MAX_ASSIGNMENT_QUERY_LEN of - true -> From + ?MAX_ASSIGNMENT_QUERY_LEN; - false -> RequestedTo - end, - { - read_local_assignments(ProcID, From, ComputedTo, Opts), - ComputedTo < RequestedTo - }. - -%% @doc Get the assignments for a process. -read_local_assignments(_ProcID, From, To, _Opts) when From > To -> - []; -read_local_assignments(ProcID, CurrentSlot, To, Opts) -> - case dev_scheduler_cache:read(ProcID, CurrentSlot, Opts) of - not_found -> - % No assignment found in cache. - []; - {ok, Assignment} -> - [ - Assignment - | read_local_assignments( - ProcID, - CurrentSlot + 1, - To, - Opts - ) - ] - end. + lib_scheduler:read_assignment_range( + dev_scheduler_cache, ProcID, From, RequestedTo, Opts). %% @doc Returns the current state of the scheduler. checkpoint(State) -> {ok, State}. diff --git a/src/preloaded/process/dev_scheduler_cache.erl b/src/preloaded/process/dev_scheduler_cache.erl index 4ec202c77..43481f4b6 100644 --- a/src/preloaded/process/dev_scheduler_cache.erl +++ b/src/preloaded/process/dev_scheduler_cache.erl @@ -10,52 +10,11 @@ %% @doc Merge the scheduler store with the main store. Used before writing %% to the cache. -opts(Opts) -> - Opts#{ - <<"store">> => - hb_opts:get( - scheduler_store, - hb_opts:get(store, no_viable_store, Opts), - Opts - ) - }. +opts(Opts) -> lib_scheduler:cache_opts(Opts). %% @doc Write an assignment message into the cache. -write(RawAssignment, RawOpts) -> - Assignment = hb_cache:ensure_all_loaded(RawAssignment, RawOpts), - Opts = opts(RawOpts), - Store = hb_opts:get(store, no_viable_store, Opts), - % Write the message into the main cache - ProcID = hb_ao:get(<<"process">>, Assignment, Opts), - Slot = hb_ao:get(<<"slot">>, Assignment, Opts), - ?event( - {writing_assignment, - {proc_id, ProcID}, - {slot, Slot}, - {assignment, Assignment} - } - ), - case hb_cache:write(Assignment, Opts) of - {ok, _UnsignedID} -> - % Create symlinks from the message on the process and the - % slot on the process to the underlying data. - ok = hb_store:link( - Store, - #{ - hb_path:to_binary([ - ?SCHEDULER_CACHE_PREFIX, - <<"assignments">>, - hb_util:human_id(ProcID), - hb_ao:normalize_key(Slot) - ]) => hb_message:id(Assignment, signed, Opts) - }, - Opts - ), - ok; - {error, Reason} -> - ?event(error, {failed_to_write_assignment, {reason, Reason}}), - {error, Reason} - end. +write(Assignment, Opts) -> + lib_scheduler:write_assignment(?SCHEDULER_CACHE_PREFIX, Assignment, Opts). %% @doc Write the initial assignment message to the cache. write_spawn(RawInitMessage, Opts) -> @@ -63,53 +22,8 @@ write_spawn(RawInitMessage, Opts) -> hb_cache:write(InitMessage, opts(Opts)). %% @doc Get an assignment message from the cache. -read(ProcID, Slot, Opts) when is_integer(Slot) -> - read(ProcID, hb_util:bin(Slot), Opts); -read(ProcID, Slot, RawOpts) -> - Opts = opts(RawOpts), - Store = hb_opts:get(store, no_viable_store, Opts), - P1 = hb_path:to_binary([ - ?SCHEDULER_CACHE_PREFIX, - <<"assignments">>, - hb_util:human_id(ProcID), - Slot - ]), - ?event( - {read_assignment, - {proc_id, ProcID}, - {slot, Slot}, - {store, Store} - } - ), - case hb_store:resolve(Store, P1, Opts) of - {ok, ResolvedPath} -> - ?event({resolved_path, {p1, P1}, {p2, ResolvedPath}, {resolved, ResolvedPath}}), - case hb_cache:read(ResolvedPath, Opts) of - {ok, RawAssignment} -> - % `hb_cache:read' no longer normalizes commitments; the - % scheduler relies on each assignment carrying its unsigned - % commitment ID, so we restore it here. - Assignment = - hb_message:normalize_commitments(RawAssignment, Opts), - % If the slot key is not present, the format of the assignment is - % AOS2, so we need to convert it to the canonical format. - case hb_ao:get(<<"variant">>, Assignment, Opts) of - <<"ao.TN.1">> -> - Loaded = hb_cache:ensure_all_loaded(Assignment, Opts), - Norm = dev_scheduler_formats:aos2_to_assignment(Loaded, Opts), - ?event({normalized_aos2_assignment, Norm}), - {ok, Norm}; - <<"ao.N.1">> -> - {ok, hb_cache:ensure_all_loaded(Assignment, Opts)} - end; - {error, not_found} -> - ?event(debug_sched, {read_assignment, {res, not_found}}), - not_found - end; - {error, not_found} -> - ?event(debug_sched, {read_assignment, {res, not_found}}), - not_found - end. +read(ProcID, Slot, Opts) -> + lib_scheduler:read_assignment(?SCHEDULER_CACHE_PREFIX, ProcID, Slot, Opts). %% @doc Get the assignments for a process. list(ProcID, RawOpts) -> diff --git a/src/preloaded/process/dev_scheduler_formats.erl b/src/preloaded/process/dev_scheduler_formats.erl index 915114f87..b12633075 100644 --- a/src/preloaded/process/dev_scheduler_formats.erl +++ b/src/preloaded/process/dev_scheduler_formats.erl @@ -238,9 +238,4 @@ aos2_normalize_types(Msg) -> %% @doc For all scheduler format operations, we do not calculate hashpaths, %% perform cache lookups, or await inprogress results. -format_opts(Opts) -> - Opts#{ - <<"hashpath">> => ignore, - <<"cache-control">> => [<<"no-cache">>, <<"no-store">>], - <<"await-inprogress">> => false - }. +format_opts(Opts) -> lib_scheduler:format_opts(Opts). diff --git a/src/preloaded/process/dev_scheduler_server.erl b/src/preloaded/process/dev_scheduler_server.erl index 3a4f4ea95..45b94438b 100644 --- a/src/preloaded/process/dev_scheduler_server.erl +++ b/src/preloaded/process/dev_scheduler_server.erl @@ -212,27 +212,23 @@ do_assign(State, Message, ReplyPID) -> BaseStateHashpath = base_state(State), NextSlot = maps:get(current, State) + 1, {Timestamp, Height, Hash} = ar_timestamp:get(), + BaseAssignment = + lib_scheduler:base_assignment( + hb_util:id(maps:get(id, State)), + NextSlot, + Message, + Opts + ), Assignment = commit_assignment( - #{ - <<"path">> => - case hb_path:from_message(request, Message, Opts) of - undefined -> <<"compute">>; - Path -> hb_path:to_binary(Path) - end, - <<"data-protocol">> => <<"ao">>, - <<"variant">> => <<"ao.N.1">>, - <<"process">> => hb_util:id(maps:get(id, State)), - <<"epoch">> => <<"0">>, - <<"slot">> => NextSlot, + BaseAssignment#{ <<"block-height">> => Height, <<"block-hash">> => hb_util:human_id(Hash), <<"block-timestamp">> => Timestamp, % Note: Local time on the SU, not Arweave <<"timestamp">> => scheduler_time(), <<"base-hashpath">> => BaseStateHashpath, - <<"body">> => OnlyAttested, - <<"type">> => <<"Assignment">> + <<"body">> => OnlyAttested }, State ), diff --git a/src/preloaded/process/lib_scheduler.erl b/src/preloaded/process/lib_scheduler.erl new file mode 100644 index 000000000..ad0fa4b83 --- /dev/null +++ b/src/preloaded/process/lib_scheduler.erl @@ -0,0 +1,322 @@ +%%% @doc A library of functions shared by the scheduler devices +%%% (`~scheduler@1.0' and `~arweave-scheduler@1.0'). Both present the same +%%% consumer interface, so the logic for interpreting a scheduling request -- +%%% locating the target process and the message to schedule -- is identical +%%% regardless of how the schedule is sourced. Reading a contiguous range of +%%% assignments from a scheduler cache is likewise shared, parameterized by +%%% the cache module. +-module(lib_scheduler). +-include("include/hb.hrl"). +-export([find_target_id/4, find_target_id/3, at_slot/2]). +-export([find_message_to_schedule/3, load_message_to_schedule/3]). +-export([only_committed/2, base_assignment/4, slot_unavailable/0]). +-export([parse_slot_range/2, read_assignment_range/5, read_local_assignments/5]). +-export([cache_opts/1, write_assignment/3, read_assignment/4, format_opts/1]). +-export([max_assignment_query_len/0]). + +%%% The maximum number of assignments that a schedule request returns at a +%%% time. +-define(MAX_ASSIGNMENT_QUERY_LEN, 1000). + +%% @doc The maximum number of assignments that a schedule request returns at a +%% time, for callers that need the limit itself (for example, to bound a +%% remote schedule request). +max_assignment_query_len() -> ?MAX_ASSIGNMENT_QUERY_LEN. + +%% @doc Find the schedule ID from a given request. The precedence order for +%% search is as follows: +%% 1. `ToSched/id' when `ToSched' has `type: Process' +%% 2. `ToSched/target' when `ToSched' has a `target' key +%% 3. `Req/target' +%% 4. `Req/id' when `Req' has `type: Process' +%% 5. `Base/process/id' +%% 6. `Base/id' when `Base' has `type: Process' +%% 7. `Req/id' +find_target_id(Base, Req, ToSched, Opts) -> + case hb_ao:get(<<"type">>, ToSched, not_found, Opts) of + <<"Process">> -> + lib_process:process_id(ToSched, #{}, Opts); + _ -> + case hb_ao:get(<<"target">>, ToSched, not_found, Opts) of + not_found -> find_target_id(Base, Req, Opts); + Target -> hb_util:human_id(Target) + end + end. +find_target_id(Base, Req, Opts) -> + TempOpts = Opts#{ <<"hashpath">> => ignore }, + Res = case hb_ao:resolve(Req, <<"target">>, TempOpts) of + {ok, Target} -> + % ID found at Req/target + Target; + _ -> + case hb_ao:resolve(Req, <<"type">>, TempOpts) of + {ok, <<"Process">>} -> + % Req is a Process, so the ID is at Req/id + lib_process:process_id(Req, #{}, Opts); + _ -> + case hb_ao:resolve(Base, <<"process">>, TempOpts) of + {ok, _Process} -> + lib_process:process_id(Base, #{}, Opts); + _ -> + % Does the message have a type of process? + case hb_ao:get(<<"type">>, Base, TempOpts) of + <<"Process">> -> + % Yes: Base is the process. + lib_process:process_id(Base, #{}, Opts); + _ -> + % No: Req is the target process. + lib_process:process_id(Req, #{}, Opts) + end + end + end + end, + ?event({found_id, {id, Res}, {base, Base}, {req, Req}}), + Res. + +%% @doc Search the given base and request message pair to find the message to +%% schedule. The precedence order for search is as follows: +%% 1. A key in `Req' with the value `self', indicating that the entire message +%% is the subject. +%% 2. A key in `Req' with another value, present in that message. +%% 3. The body of the message. +%% 4. The message itself. +find_message_to_schedule(Base, Req, Opts) -> + Subject = + hb_ao:get( + <<"subject">>, + Req, + not_found, + Opts#{ <<"hashpath">> => ignore } + ), + case Subject of + <<"base">> -> Base; + <<"self">> -> Req; + not_found -> + hb_ao:get(<<"body">>, Req, Req, Opts#{ <<"hashpath">> => ignore }); + Subject -> + hb_ao:get(Subject, Req, Opts#{ <<"hashpath">> => ignore }) + end. + +%% @doc The slot that the given (process-shaped) base message has been +%% computed to. +at_slot(Base, Opts) -> + hb_util:int( + hb_ao:get( + <<"at-slot">>, + Base, + Opts#{ <<"hashpath">> => ignore } + ) + ). + +%% @doc Reduce a message to its committed components, returning a 400 error +%% message if they cannot be validated. +only_committed(Msg, Opts) -> + case hb_message:with_only_committed(Msg, Opts) of + {ok, OnlyCommitted} -> {ok, OnlyCommitted}; + {error, Err} -> + {error, + #{ + <<"status">> => 400, + <<"body">> => <<"Message invalid: ", + "Committed components cannot be validated.">>, + <<"reason">> => Err + } + } + end. + +%% @doc The keys of an assignment that are shared by every scheduler device: +%% everything except the fields that describe the assignment's position in the +%% scheduler's specific sequencing space (block info, weave offset, et al), +%% which the caller merges in. The `path' is taken from `PathMsg' -- normally +%% the assigned message itself. +base_assignment(ProcID, Slot, PathMsg, Opts) -> + #{ + <<"path">> => + case hb_path:from_message(request, PathMsg, Opts) of + undefined -> <<"compute">>; + Path -> hb_path:to_binary(Path) + end, + <<"data-protocol">> => <<"ao">>, + <<"variant">> => <<"ao.N.1">>, + <<"process">> => ProcID, + <<"epoch">> => <<"0">>, + <<"slot">> => Slot, + <<"body">> => PathMsg, + <<"type">> => <<"Assignment">> + }. + +%% @doc The error returned when a requested slot is not yet present in the +%% schedule. +slot_unavailable() -> + {error, + #{ + <<"status">> => 404, + <<"reason">> => + <<"Requested slot not yet available in schedule.">> + } + }. + +%% @doc Find and fully load the message to schedule from the given base and +%% request pair. If the message cannot be completely loaded, a 404 error +%% message is returned instead. +load_message_to_schedule(Base, Req, Opts) -> + RawToSched = find_message_to_schedule(Base, Req, Opts), + try {ok, hb_cache:ensure_all_loaded(RawToSched, Opts)} + catch + error:{necessary_message_not_found, _, _} -> + {error, + #{ + <<"status">> => 404, + <<"body">> => <<"Cannot fully load message to schedule.">> + } + } + end. + +%% @doc Parse the requested slot range -- `from' and `to' -- from a schedule +%% request. `from' defaults to 0 and is clamped to be non-negative; `to' +%% defaults to `undefined' (the latest slot known to the caller). +parse_slot_range(Req, Opts) -> + From = + case hb_ao:get(<<"from">>, Req, not_found, Opts) of + not_found -> 0; + FromRes -> max(0, hb_util:int(FromRes)) + end, + To = + case hb_ao:get(<<"to">>, Req, not_found, Opts) of + not_found -> undefined; + ToRes -> hb_util:int(ToRes) + end, + {From, To}. + +%% @doc Read the assignments for slots `From'..`RequestedTo' from a scheduler +%% cache, truncating the range to at most `?MAX_ASSIGNMENT_QUERY_LEN' slots. +%% Returns the assignments and whether the request was truncated. +read_assignment_range(CacheMod, ProcID, From, RequestedTo, Opts) -> + ComputedTo = min(RequestedTo, From + ?MAX_ASSIGNMENT_QUERY_LEN), + { + read_local_assignments(CacheMod, ProcID, From, ComputedTo, Opts), + ComputedTo < RequestedTo + }. + +%% @doc Read a contiguous range of assignments (slots `From'..`To', inclusive) +%% from a scheduler cache, stopping at the first slot that is not present. The +%% cache module is passed in so both scheduler devices can share the traversal: +%% `CacheMod:read(ProcID, Slot, Opts)' must return `{ok, Assignment}' or +%% `not_found'. +read_local_assignments(_CacheMod, _ProcID, From, To, _Opts) when From > To -> + []; +read_local_assignments(CacheMod, ProcID, From, To, Opts) -> + case CacheMod:read(ProcID, From, Opts) of + not_found -> []; + {ok, Assignment} -> + [ + Assignment + | read_local_assignments(CacheMod, ProcID, From + 1, To, Opts) + ] + end. + +%% @doc Options for generating (and reading messages for) schedule responses: +%% formatting is deterministic, so its resolutions are neither cached nor +%% awaited. +format_opts(Opts) -> + Opts#{ + <<"hashpath">> => ignore, + <<"cache-control">> => [<<"no-cache">>, <<"no-store">>], + <<"await-inprogress">> => false + }. + +%% @doc Merge the scheduler store with the main store. Used before reading +%% from or writing to a scheduler cache. +cache_opts(Opts) -> + Opts#{ + <<"store">> => + hb_opts:get( + scheduler_store, + hb_opts:get(store, no_viable_store, Opts), + Opts + ) + }. + +%% @doc Write an assignment message into a scheduler cache: the message goes +%% into the main cache, and `/assignments//' is linked to +%% its signed ID. The pseudo-path prefix distinguishes the caches of the +%% scheduler devices that share this helper. +write_assignment(Prefix, RawAssignment, RawOpts) -> + Assignment = hb_cache:ensure_all_loaded(RawAssignment, RawOpts), + Opts = cache_opts(RawOpts), + Store = hb_opts:get(store, no_viable_store, Opts), + ProcID = hb_ao:get(<<"process">>, Assignment, Opts), + Slot = hb_ao:get(<<"slot">>, Assignment, Opts), + ?event( + {writing_assignment, + {prefix, Prefix}, + {proc_id, ProcID}, + {slot, Slot} + } + ), + case hb_cache:write(Assignment, Opts) of + {ok, _UnsignedID} -> + ok = hb_store:link( + Store, + #{ + assignment_path(Prefix, ProcID, Slot) => + hb_message:id(Assignment, signed, Opts) + }, + Opts + ), + ok; + {error, Reason} -> + ?event(error, {failed_to_write_assignment, {reason, Reason}}), + {error, Reason} + end. + +%% @doc Get an assignment message from a scheduler cache. Restores the +%% assignment's unsigned commitment ID (`hb_cache:read' does not normalize +%% commitments) and converts legacy `ao.TN.1' (AOS2) assignments to the +%% canonical format via the scheduler package's formats module. +read_assignment(Prefix, ProcID, Slot, Opts) when is_integer(Slot) -> + read_assignment(Prefix, ProcID, hb_util:bin(Slot), Opts); +read_assignment(Prefix, ProcID, Slot, RawOpts) -> + Opts = cache_opts(RawOpts), + Store = hb_opts:get(store, no_viable_store, Opts), + Path = assignment_path(Prefix, ProcID, Slot), + ?event( + {read_assignment, + {prefix, Prefix}, + {proc_id, ProcID}, + {slot, Slot} + } + ), + case hb_store:resolve(Store, Path, Opts) of + {ok, ResolvedPath} -> + case hb_cache:read(ResolvedPath, Opts) of + {ok, RawAssignment} -> + Assignment = + hb_message:normalize_commitments(RawAssignment, Opts), + case hb_ao:get(<<"variant">>, Assignment, Opts) of + <<"ao.TN.1">> -> + Loaded = + hb_cache:ensure_all_loaded(Assignment, Opts), + {ok, + dev_scheduler_formats:aos2_to_assignment( + Loaded, + Opts + ) + }; + <<"ao.N.1">> -> + {ok, hb_cache:ensure_all_loaded(Assignment, Opts)} + end; + {error, not_found} -> not_found + end; + {error, not_found} -> not_found + end. + +assignment_path(Prefix, ProcID, Slot) -> + hb_path:to_binary( + [ + Prefix, + <<"assignments">>, + hb_util:human_id(ProcID), + hb_ao:normalize_key(Slot) + ] + ). From 0364fd3174ec5ba238d1808e8fd51a15738b36b7 Mon Sep 17 00:00:00 2001 From: Sam Williams Date: Thu, 16 Jul 2026 14:11:27 -0400 Subject: [PATCH 3/7] refactor(arweave): share best_response via lib_arweave_common `dev_arweave' and the `~arweave-scheduler@1.0' device both broadcast a request to several Arweave nodes and take the most successful reply. Move `dev_arweave''s status-sorting `best_response/1' (with its `response_status/1' helper) into `lib_arweave_common', which both devices already depend on, and drop the local copy. No behavioural change for `dev_arweave'. --- src/preloaded/arweave/dev_arweave.erl | 28 ++-------------------- src/preloaded/codec/lib_arweave_common.erl | 27 +++++++++++++++++++++ 2 files changed, 29 insertions(+), 26 deletions(-) diff --git a/src/preloaded/arweave/dev_arweave.erl b/src/preloaded/arweave/dev_arweave.erl index 580b069cf..7dcec5ece 100644 --- a/src/preloaded/arweave/dev_arweave.erl +++ b/src/preloaded/arweave/dev_arweave.erl @@ -885,31 +885,7 @@ request(Method, Path, Extra, LogExtra, Opts) -> <<"cache-control">> => [<<"no-cache">>, <<"no-store">>] } ), - to_message(Path, Method, best_response(Res), LogExtra, Opts). - -%% @doc Select the best response from a list of responses by sorting them -%% ascending by HTTP status code. Returns the first (best) response tuple. -best_response({error, {no_viable_responses, Responses}}) -> - best_response(Responses); -best_response([]) -> - {error, no_viable_responses}; -best_response(Responses) when is_list(Responses) -> - Sorted = lists:sort( - fun({_, ResponseA}, {_, ResponseB}) -> - StatusA = response_status(ResponseA), - StatusB = response_status(ResponseB), - StatusA =< StatusB - end, - Responses - ), - hd(Sorted); -best_response(Response) -> - Response. - -response_status(Response) when is_map(Response) -> - maps:get(<<"status">>, Response, 999); -response_status(_Response) -> - 999. + to_message(Path, Method, lib_arweave_common:best_response(Res), LogExtra, Opts). %% @doc Transform a response from the Arweave node into an AO-Core message. to_message(Path, Method, {error, #{ <<"status">> := 404 }}, LogExtra, _Opts) -> @@ -1206,7 +1182,7 @@ best_response_handles_failed_connect_entries_test_parallel() -> ], ?assertEqual( {ok, #{ <<"status">> => 200, <<"body">> => <<"OK-2">> }}, - best_response(Responses) + lib_arweave_common:best_response(Responses) ). best_response_non_map_error_round_trips_test_parallel() -> diff --git a/src/preloaded/codec/lib_arweave_common.erl b/src/preloaded/codec/lib_arweave_common.erl index 755b9492c..6430ad2de 100644 --- a/src/preloaded/codec/lib_arweave_common.erl +++ b/src/preloaded/codec/lib_arweave_common.erl @@ -6,6 +6,7 @@ -export([bundle_hint/4, data/3, tags/5, excluded_tags/3]). -export([to/3, to/6, siginfo/4, fields_to_tx/4]). -export([bundle_header/2, bundle_header/3]). +-export([best_response/1]). -include("include/hb.hrl"). -define(ANS104_BASE_FIELDS, [<<"anchor">>, <<"target">>]). @@ -770,3 +771,29 @@ read_bundle_header(BundleStartOffset, HeaderSize, FirstChunk, Opts) -> Error -> Error end. + +%% @doc Select the best response from a (potentially multi-node) request result +%% by sorting the responses ascending by HTTP status code and returning the +%% first (best) one. Shared by the Arweave devices, which broadcast requests to +%% several nodes and take the most successful reply. +best_response({error, {no_viable_responses, Responses}}) -> + best_response(Responses); +best_response([]) -> + {error, no_viable_responses}; +best_response(Responses) when is_list(Responses) -> + Sorted = lists:sort( + fun({_, ResponseA}, {_, ResponseB}) -> + StatusA = response_status(ResponseA), + StatusB = response_status(ResponseB), + StatusA =< StatusB + end, + Responses + ), + hd(Sorted); +best_response(Response) -> + Response. + +response_status(Response) when is_map(Response) -> + maps:get(<<"status">>, Response, 999); +response_status(_Response) -> + 999. From c60a91002c1454f16ab5692799d05b5963d8454a Mon Sep 17 00:00:00 2001 From: Sam Williams Date: Thu, 16 Jul 2026 14:11:27 -0400 Subject: [PATCH 4/7] fix(query): resolve `block` on ~query@1.0 transaction nodes A transaction node returned from the local ~query@1.0 GraphQL endpoint resolved `block' to null, because a cached tx carries no block metadata of its own. Carry each match's weave offset onto its node and resolve `block' by locating the cached block whose byte range contains that offset, distinguished from the top-level block(id/height) query by the offset key on the node. Resolution is local but scans the cached blocks, so callers that only need on-chain ordering should use the offset directly. The cached-heights enumeration and the block byte-bounds arithmetic are shared as `cached_heights/1' and `block_bounds/2' with the module's pre-existing `latest_cached_block/1' and `block_range_to_offset_range/2', which previously inlined the same code separately. --- src/preloaded/query/dev_query_arweave.erl | 80 +++++++++++++++++------ 1 file changed, 60 insertions(+), 20 deletions(-) diff --git a/src/preloaded/query/dev_query_arweave.erl b/src/preloaded/query/dev_query_arweave.erl index 5ffb05f3c..6b91e847d 100644 --- a/src/preloaded/query/dev_query_arweave.erl +++ b/src/preloaded/query/dev_query_arweave.erl @@ -87,6 +87,13 @@ query(Obj, <<"transactions">>, Args, Opts) -> ), {ok, connection([], Args, Opts)} end; +query(#{ <<"offset">> := Offset }, <<"block">>, _Args, Opts) + when is_integer(Offset) -> + % Resolve the `block' field of a transaction node: a cached transaction + % holds no block metadata of its own, so the block is found by the weave + % offset carried on the node (see `node_with_offset/3'), which also + % distinguishes this from the top-level `block(id/height)' query. + block_at_offset(Offset, Opts); query(Obj, <<"block">>, Args, Opts) -> case query(Obj, <<"blocks">>, Args, Opts) of {ok, []} -> {ok, null}; @@ -264,11 +271,21 @@ read_ids(_, 0, _Opts) -> []; read_ids([AnnotatedID = #{ <<"id">> := ID } | Rest], Count, Opts) -> case hb_cache:read(ID, Opts) of {ok, Msg} -> - [AnnotatedID#{ <<"node">> => Msg } | read_ids(Rest, Count - 1, Opts)]; + [ + AnnotatedID#{ <<"node">> => node_with_offset(Msg, AnnotatedID, Opts) } + | read_ids(Rest, Count - 1, Opts) + ]; _ -> read_ids(Rest, Count, Opts) end. +%% @doc Carry the transaction's weave offset onto its node, so that the `block' +%% field resolver can locate the block that includes it (see `query/4'). +node_with_offset(Msg, #{ <<"offset">> := Offset }, _Opts) + when is_map(Msg), is_integer(Offset) -> + Msg#{ <<"offset">> => Offset }; +node_with_offset(Msg, _Annotated, _Opts) -> Msg. + %% @doc Drop to the cursor position, returning the list of items after the cursor. drop_to_cursor(Args, Ordered, Opts) -> drop_to_cursor( @@ -380,14 +397,8 @@ block_range_to_offset_range(Heights, Opts) -> RawMin -> case read_block(hb_util:int(RawMin), Opts) of {ok, MinBlock} -> - % The `weave_size` is the size at the _end_ of the block, - % so we must subtract the start from it to find the - % starting byte of the block. - WeaveSize = hb_util:int( - hb_maps:get(<<"weave_size">>, MinBlock, 0, Opts)), - BlockSize = hb_util:int( - hb_maps:get(<<"block_size">>, MinBlock, 0, Opts)), - WeaveSize - BlockSize; + {BlockStart, _} = block_bounds(MinBlock, Opts), + BlockStart; {error, not_found} -> 0 end end, @@ -456,18 +467,47 @@ read_cached_block(Height, Opts) -> %% @doc Return the latest block height indexed in the Arweave pseudo-path cache. latest_cached_block(Opts) -> - Blocks = - hb_cache:list_numbered( - hb_path:to_binary([ - <<"~arweave@2.9">>, - <<"block">>, - <<"height">> - ]), - Opts - ), - case Blocks of + case cached_heights(Opts) of [] -> not_found; - _ -> {ok, lists:max(Blocks)} + Blocks -> {ok, lists:max(Blocks)} + end. + +%% @doc List the block heights present in the Arweave pseudo-path cache. +cached_heights(Opts) -> + hb_cache:list_numbered( + hb_path:to_binary([ + <<"~arweave@2.9">>, + <<"block">>, + <<"height">> + ]), + Opts + ). + +%% @doc The weave byte range `{StartOffset, EndOffset}' that a block covers. +%% The block's `weave_size' is the size of the weave at the block's _end_, so +%% its starting byte is that size less the block's own size. +block_bounds(Block, Opts) -> + WeaveSize = hb_util:int(hb_maps:get(<<"weave_size">>, Block, 0, Opts)), + BlockSize = hb_util:int(hb_maps:get(<<"block_size">>, Block, 0, Opts)), + {WeaveSize - BlockSize, WeaveSize}. + +%% @doc Find the cached block that includes the given weave offset -- the block +%% whose byte range `[weave_size - block_size, weave_size)' contains it. Scans +%% the locally cached blocks, so it is only reached when the `block' field is +%% explicitly selected on a transaction. Returns `{ok, null}' when no cached +%% block covers the offset. +block_at_offset(Offset, Opts) -> + block_covering(Offset, lists:sort(cached_heights(Opts)), Opts). + +block_covering(_Offset, [], _Opts) -> {ok, null}; +block_covering(Offset, [Height | Rest], Opts) -> + case read_cached_block(Height, Opts) of + {ok, Block} -> + case block_bounds(Block, Opts) of + {Start, End} when Start =< Offset, Offset < End -> {ok, Block}; + _ -> block_covering(Offset, Rest, Opts) + end; + _ -> block_covering(Offset, Rest, Opts) end. %%% Match argument processing From 97ea6c2db6e567609177f4c726d909f471696062 Mon Sep 17 00:00:00 2001 From: Sam Williams Date: Thu, 16 Jul 2026 14:11:47 -0400 Subject: [PATCH 5/7] feat: add ~arweave-scheduler@1.0, an Arweave-L1-sequenced scheduler A drop-in ~scheduler@1.0 replacement that sources a process's schedule directly from Arweave layer-1 transactions. Slot 0 is the process message (its tx id is the process id); slot N is the Nth confirmed L1 tx whose `target' is the process id, in ascending weave-offset order. Assignments are deterministic and uncommitted -- Arweave itself is the authority -- so every node converges on identical assignment ids without trusting a scheduler wallet. The on-chain position recorded on each assignment is the weave offset. Discovery is local-first, with no gateway attribution: a single `~copycat@1.0/arweave' shallow pass over the range records each tx's weave offset and caches its (data-free) header, which populates the local `field-target' match; the device then resolves the node's own ~query@1.0 (via hb_ao:resolve, not a remote call) for the recipient set, ordered by that offset index. `arweave_scheduler_query_source => remote' opts into a gateway query instead. The block index is node-shared and `ensure_offsets' passes `reindex=false', so N processes on one node index each block once rather than once per process. The device is header-only end to end. Bodies are read as tx headers, so the schedule never depends on data availability; the header verifies and its committed id is the tx id. POST /schedule likewise relays the header only, rejecting a message that carries data (422) rather than uploading chunks, and rejecting any non-tx@1.0 commitment (422). The misleading tip-derived timestamp/block-height/block-hash that the standard formatter stamps are omitted from a chain-derived schedule. Synchronization tracks a `synced-to' block high-water mark: the attestable range `[spawn-height, synced-to]' the node can vouch for. A read with no new blocks does no network work; otherwise only the delta is indexed, in `arweave_scheduler_sync_chunk'-block passes whose progress is persisted per chunk (resumable first sync). /status reports each tracked process's range. The device is built on `lib_scheduler': request parsing, the assignment cache, the range read, and the base assignment schema are the same code paths ~scheduler@1.0 uses. The full consumer interface (status/next/schedule/slot/init/checkpoint + router) is implemented so ~process@1.0 drives it unchanged. Includes a permanent mainnet test fixture (a lua@5.3a process whose three scheduled messages transform state 1000 -> 1337, one posted directly to arweave.net to prove foreign-tx indexing) and device unit tests. --- .../process/dev_arweave_scheduler.erl | 1042 +++++++++++++++++ .../process/dev_arweave_scheduler_cache.erl | 111 ++ test/arweave-scheduler-test.lua | 34 + 3 files changed, 1187 insertions(+) create mode 100644 src/preloaded/process/dev_arweave_scheduler.erl create mode 100644 src/preloaded/process/dev_arweave_scheduler_cache.erl create mode 100644 test/arweave-scheduler-test.lua diff --git a/src/preloaded/process/dev_arweave_scheduler.erl b/src/preloaded/process/dev_arweave_scheduler.erl new file mode 100644 index 000000000..37a32ea45 --- /dev/null +++ b/src/preloaded/process/dev_arweave_scheduler.erl @@ -0,0 +1,1042 @@ +%%% @doc A scheduler for AO processes that uses the Arweave network itself +%%% as the sequencing layer. The device implements the same interface as +%%% `~scheduler@1.0' from the perspective of consumers (`~process@1.0' et +%%% al), but assignments are not minted by a scheduling authority: they are +%%% implicit in the order that Arweave blocks include L1 transactions that +%%% target the process. +%%% +%%% The model is as follows: +%%%
    +%%%
  • A process is spawned by uploading its process message to Arweave +%%% as a `tx@1.0'-committed L1 transaction. The transaction ID is the +%%% process ID, and the process message occupies slot 0 of its own +%%% schedule.
  • +%%%
  • Every L1 transaction whose `target' field is the process ID is a +%%% message in the process's schedule. Slots follow the canonical weave +%%% order (ascending weave offset, which is ascending block order). The +%%% transaction's weave `offset' -- not a scheduler-assigned nonce -- is +%%% its on-chain position, and it is recorded on the assignment.
  • +%%%
  • Discovery is local-first: the node indexes the relevant +%%% blocks itself with `~copycat@1.0' and then queries its own +%%% `~query@1.0' GraphQL endpoint (`transactions(recipients: +%%% [ProcessID], sort: HEIGHT_ASC)') for the base-layer transactions +%%% addressed to the process, ordered by the local weave-offset index. +%%% A node may instead be configured to query a remote gateway +%%% (`arweave_scheduler_query_source => remote'). There is no bespoke +%%% store index -- discovery uses the node's existing Arweave index and +%%% the device's own schedule cache.
  • +%%%
+%%% +%%% Assignment bodies are transaction headers only: each is the +%%% data-free header that `~copycat@1.0/arweave' cached while indexing the +%%% range (read from the node's stores), so the schedule never carries (nor +%%% depends on the availability of) a transaction's data. A header still +%%% carries the tx@1.0 signature, so its committed ID is the transaction ID +%%% and it verifies independently. Assignments are +%%% deterministic derivations of chain data and are left uncommitted: every +%%% node that reads the same blocks converges on an identical schedule +%%% without trusting a scheduler wallet. Only `tx@1.0' commitments are +%%% accepted for dispatch -- ANS-104 data items and HTTPSig messages cannot +%%% be sequenced by the base layer. +%%% +%%% `POST /schedule' dispatches a user's presigned transaction to the +%%% network on their behalf. Only its header is relayed: because the schedule +%%% is built from headers, a message carrying data is rejected rather than +%%% dispatched, so its slot can never depend on chunk availability. No slot is +%%% returned at dispatch time: the transaction receives its slot once it is +%%% included in a block. Note +%%% that Arweave refuses to create accounts with a zero balance +%%% (`validate_overspend'), and a process ID is a fresh account: the first +%%% message addressed to a process must therefore carry a `quantity' of at +%%% least 1 winston (and, until the account exists, the sender's `reward' +%%% must cover the network's new-account fee). +%%% +%%% Each process's synchronization tracks the contiguous block range +%%% `[spawn-height, synced-to]' it has indexed -- the range the node can +%%% attest to. A sync only indexes the blocks above `synced-to' (up to the +%%% confirmed tip), so a repeated read with no new blocks does no network +%%% work, and `/status' reports each tracked process's range. The device +%%% honors the following node options: +%%%
    +%%%
  • `arweave_scheduler_confirmation_depth': blocks under the network +%%% tip that must elapse before a block is considered final +%%% (default: 10).
  • +%%%
  • `arweave_scheduler_sync_chunk': blocks indexed per synchronization +%%% pass, bounding the work redone if a pass is interrupted +%%% (default: 1000).
  • +%%%
  • `arweave_scheduler_max_height': optional hard upper bound on the +%%% synced height, pinning a schedule to an immutable block range +%%% (used primarily in tests).
  • +%%%
+-module(dev_arweave_scheduler). +-implements(<<"arweave-scheduler@1.0">>). +-device_libraries([lib_process, lib_scheduler, lib_arweave_common]). +%%% AO-Core API functions: +-export([info/0]). +%%% Scheduling functions: +-export([schedule/3, router/4]). +%%% CU-flow functions: +-export([slot/3, status/3, next/3]). +-export([checkpoint/1]). +-include("include/hb.hrl"). +-include_lib("eunit/include/eunit.hrl"). + +%%% The default number of blocks beneath the network tip at which we +%%% consider a block final. +-define(DEFAULT_CONFIRMATION_DEPTH, 10). +%%% The default number of blocks indexed per synchronization pass. A sync +%%% persists its progress after each chunk, so this bounds the work redone if +%%% a pass is interrupted and keeps the first sync of a long-lived process +%%% resumable. Overridable with the `arweave_scheduler_sync_chunk' option. +-define(DEFAULT_SYNC_CHUNK_BLOCKS, 1000). +%%% The number of transactions requested per query page. +-define(QUERY_PAGE_SIZE, 100). +%%% The Arweave device that we resolve chain data through. +-define(ARWEAVE_DEVICE, <<"~arweave@2.9">>). + +%% @doc This device uses a default handler to route requests to the correct +%% function. +info() -> + #{ + exports => + [ + <<"status">>, + <<"next">>, + <<"schedule">>, + <<"slot">>, + <<"init">>, + <<"checkpoint">> + ], + excludes => [set, keys], + default => fun router/4 + }. + +%% @doc The default handler for the device: route all unrecognized requests +%% to `schedule'. +router(_, Base, Req, Opts) -> + ?event({arweave_scheduler_router_called, {req, Req}}), + schedule(Base, Req, Opts). + +%% @doc Return the next assignment for a process. Assumes that `Base' is a +%% `dev_process' or similar message, having an `at-slot' key. If the next +%% slot is not present in the local cache, the schedule is synchronized +%% from Arweave before the read is retried. +next(Base, Req, Opts) -> + ProcID = lib_process:process_id(Base, Req, Opts), + LastProcessed = lib_scheduler:at_slot(Base, Opts), + ?event(next, {arweave_next, {proc_id, ProcID}, {last, LastProcessed}}), + maybe + {ok, Assignment} ?= find_assignment(ProcID, LastProcessed + 1, Opts), + {ok, #{ <<"body">> => Assignment, <<"state">> => Base }} + end. + +%% @doc Read an assignment from the local cache, synchronizing the schedule +%% from Arweave if it is not yet present. +find_assignment(ProcID, Slot, Opts) -> + case dev_arweave_scheduler_cache:read(ProcID, Slot, Opts) of + {ok, Assignment} -> {ok, Assignment}; + not_found -> + maybe + {ok, _} ?= sync(ProcID, Opts), + case dev_arweave_scheduler_cache:read(ProcID, Slot, Opts) of + {ok, Assignment} -> {ok, Assignment}; + not_found -> lib_scheduler:slot_unavailable() + end + end + end. + +%% @doc A router for choosing between getting the existing schedule, or +%% dispatching a new message to Arweave. +schedule(Base, Req, Opts) -> + ?event({resolving_arweave_schedule_request, {req, Req}}), + case hb_util:key_to_atom(hb_ao:get(<<"method">>, Req, <<"GET">>, Opts)) of + post -> post_schedule(Base, Req, Opts); + get -> get_schedule(Base, Req, Opts) + end. + +%% @doc Generate and return the schedule for a process, optionally between +%% two slots -- labelled as `from' and `to'. +get_schedule(Base, Req, Opts) -> + ProcID = hb_util:human_id(lib_scheduler:find_target_id(Base, Req, Opts)), + {From, To} = lib_scheduler:parse_slot_range(Req, Opts), + maybe + {ok, #{ <<"next-slot">> := NextSlot }} ?= sync(ProcID, Opts), + Latest = NextSlot - 1, + RequestedTo = + case To of + undefined -> Latest; + _ -> min(To, Latest) + end, + {Assignments, More} = + lib_scheduler:read_assignment_range( + dev_arweave_scheduler_cache, ProcID, From, RequestedTo, Opts), + dev_arweave_scheduler_cache:assignments_to_bundle( + ProcID, + Assignments, + More, + Opts + ) + end. + +%% @doc Dispatch a new message to Arweave. The message must carry a signed +%% `tx@1.0' commitment: the caller signs the L1 transaction themselves +%% (including its `anchor' and `reward'), and this device relays it to the +%% network. The message is also written to the local cache, so that it is +%% servable by ID before the network has propagated and indexed it. +post_schedule(Base, Req, Opts) -> + maybe + {ok, ToSched} ?= lib_scheduler:load_message_to_schedule(Base, Req, Opts), + do_post_schedule(Base, Req, ToSched, Opts) + end. + +do_post_schedule(Base, Req, ToSched, Opts) -> + ProcID = lib_scheduler:find_target_id(Base, Req, ToSched, Opts), + ?event({arweave_post_schedule, {proc_id, ProcID}}), + maybe + {ok, OnlyCommitted} ?= lib_scheduler:only_committed(ToSched, Opts), + ok ?= ensure_tx_committed(OnlyCommitted, Opts), + dispatch(ProcID, OnlyCommitted, Opts) + end. + +%% @doc Ensure that the given message carries a signed, valid `tx@1.0' +%% commitment. The Arweave scheduler accepts only L1 transactions: ANS-104 +%% and HTTPSig commitments cannot be sequenced by the base layer. +ensure_tx_committed(Msg, Opts) -> + Devices = hb_message:commitment_devices(Msg, Opts), + Signers = hb_message:signers(Msg, Opts), + case {lists:member(<<"tx@1.0">>, Devices), Signers} of + {false, _} -> + {error, + #{ + <<"status">> => 422, + <<"body">> => + <<"The Arweave scheduler only accepts messages ", + "committed with tx@1.0.">> + } + }; + {true, []} -> + {error, + #{ + <<"status">> => 422, + <<"body">> => <<"Message must be signed.">> + } + }; + {true, _} -> + case hb_message:verify(Msg, signers, Opts) of + true -> ok; + false -> + {error, + #{ + <<"status">> => 400, + <<"body">> => <<"Message is not valid.">> + } + } + end + end. + +%% @doc Relay a committed `tx@1.0' message to the Arweave network. Only the +%% transaction header is dispatched: this scheduler sequences headers, never +%% data. A message carrying data is therefore rejected rather than uploaded -- +%% its chunks might be unavailable when the schedule is later read. The same +%% header-only rule governs synchronization, so a data-carrying transaction +%% posted directly to Arweave is still sequenced safely, by its header alone. +dispatch(ProcID, Msg, Opts) -> + TX = hb_message:convert(Msg, <<"tx@1.0">>, Opts), + case TX#tx.data of + <<>> -> dispatch_header(ProcID, TX, Msg, Opts); + _ -> + {error, + #{ + <<"status">> => 422, + <<"body">> => + <<"The Arweave scheduler sequences transaction ", + "headers only; a message carrying data cannot ", + "be scheduled.">> + } + } + end. + +dispatch_header(ProcID, TX, Msg, Opts) -> + TXID = hb_util:human_id(TX#tx.id), + ?event({dispatching_tx, {id, {explicit, TXID}}, {proc_id, ProcID}}), + maybe + {ok, _} ?= post_tx_header(TX, Opts), + {ok, _} = hb_cache:write(Msg, Opts), + % `tx-id' rather than `id': the reserved `id' key would be shadowed + % by the receipt message's own ID when read by consumers. + {ok, + #{ + <<"status">> => 202, + <<"tx-id">> => TXID, + <<"process">> => hb_util:human_id(ProcID), + <<"body">> => + <<"Transaction dispatched to Arweave. It will receive ", + "a slot in the schedule once mined.">> + } + } + end. + +%% @doc Post a transaction header to the network, via the node's `/arweave' +%% route. The route may broadcast to multiple nodes: acceptance by any one +%% of them places the transaction in the network's mempool, so the best +%% response wins, mirroring `~arweave@2.9's handling. +post_tx_header(TX, Opts) -> + Res = + hb_http:request( + #{ + <<"method">> => <<"POST">>, + <<"path">> => <<"/arweave/tx">>, + <<"body">> => hb_json:encode(ar_tx:tx_to_json_struct(TX)) + }, + no_result_cache(Opts) + ), + lib_arweave_common:best_response(Res). + +%% @doc Returns information about the current slot for a process. +slot(Base, Req, Opts) -> + ProcID = hb_util:human_id(lib_scheduler:find_target_id(Base, Req, Opts)), + ?event({getting_current_slot, {proc_id, ProcID}}), + maybe + {ok, #{ <<"next-slot">> := NextSlot }} ?= sync(ProcID, Opts), + {ok, + #{ + <<"process">> => ProcID, + <<"current">> => NextSlot - 1, + <<"cache-control">> => <<"no-store">> + } + } + end. + +%% @doc Returns information about the scheduler, including -- for each process +%% it tracks -- the contiguous block range it has indexed and can attest to. +status(_Base, _Req, Opts) -> + Wallet = hb_opts:get(priv_wallet, hb:wallet(), Opts), + ProcIDs = dev_arweave_scheduler_cache:list_processes(Opts), + {ok, + #{ + <<"address">> => hb_util:human_id(ar_wallet:to_address(Wallet)), + <<"processes">> => + hb_maps:from_list([ {P, attestable_range(P, Opts)} || P <- ProcIDs ]), + <<"cache-control">> => <<"no-store">> + } + }. + +%% @doc The block range a process has been indexed over, as it currently stands +%% in the cache (no synchronization is triggered). `from' is the spawn block and +%% `to' is the confirmed height up to which the schedule is complete; `current' +%% is the latest materialized slot. +attestable_range(ProcID, Opts) -> + case dev_arweave_scheduler_cache:read_state(ProcID, Opts) of + {ok, #{ + <<"spawn-height">> := Spawn, + <<"synced-to">> := SyncedTo, + <<"next-slot">> := NextSlot + }} -> + #{ + <<"from">> => Spawn, + <<"to">> => SyncedTo, + <<"current">> => NextSlot - 1 + }; + _ -> #{} + end. + +%% @doc Returns the current state of the scheduler. +checkpoint(State) -> {ok, State}. + +%%% Synchronization + +%% @doc Synchronize the schedule of a process from Arweave, extending the +%% contiguously-indexed block range `[spawn-height, synced-to]' up to the +%% confirmed network tip. Only the blocks above `synced-to' are indexed; a +%% process already synced to the tip does no network work at all. +sync(ProcID, Opts) -> + maybe + {ok, State} ?= ensure_initialized(ProcID, Opts), + {ok, Upper} ?= confirmed_tip(Opts), + do_sync(ProcID, State, Upper, Opts) + end. + +%% @doc Extend the synced range up to `Upper', one bounded chunk at a time. The +%% range grows contiguously from `synced-to + 1', and `{next-slot, synced-to}' +%% are persisted together after each chunk, so an interrupted sync resumes at +%% the start of the unfinished chunk (re-materializing is idempotent) rather +%% than re-deriving the whole history. +do_sync(_ProcID, State = #{ <<"synced-to">> := SyncedTo }, Upper, _Opts) + when SyncedTo >= Upper -> + {ok, State}; +do_sync(ProcID, State = #{ <<"synced-to">> := SyncedTo }, Upper, Opts) -> + From = SyncedTo + 1, + To = min(SyncedTo + sync_chunk_blocks(Opts), Upper), + ?event( + {arweave_scheduler_sync, + {proc_id, ProcID}, + {from, From}, + {to, To}, + {target, Upper} + } + ), + maybe + {ok, Ordered} ?= discover(ProcID, From, To, Opts), + {ok, NewState} ?= materialize(ProcID, State, Ordered, To, Opts), + do_sync(ProcID, NewState, Upper, Opts) + end. + +%% @doc The number of blocks to index per synchronization pass. +sync_chunk_blocks(Opts) -> + hb_util:int( + hb_opts:get(arweave_scheduler_sync_chunk, ?DEFAULT_SYNC_CHUNK_BLOCKS, Opts) + ). + +%% @doc Side-effecting and mutable internal resolutions must never be served +%% from (or written to) the resolution cache: nodes default to caching every +%% HTTP resolution, which would otherwise freeze the indexing runs, queries +%% and tip lookups at their first result. +no_result_cache(Opts) -> + Opts#{ + <<"hashpath">> => ignore, + <<"cache-control">> => [<<"no-cache">>, <<"no-store">>] + }. + +%% @doc Discover the base-layer transactions addressed to a process within a +%% block range, in ascending weave order, entirely from the node's own index. +%% Indexing the range with `~copycat@1.0/arweave' both records each +%% transaction's weave offset and caches its header (so its `target' is locally +%% matchable); the recipient match is then served from the node's own +%% `~query@1.0' GraphQL endpoint, and each match is annotated with its offset -- +%% the sort key and the `offset' recorded on the assignment. No gateway is +%% queried by default. Returns `{Offset, TXID}' pairs in ascending offset order; +%% bundled data items are excluded, as this scheduler sequences the base layer +%% only. +discover(ProcID, From, To, Opts) -> + maybe + ok ?= ensure_offsets(From, To, Opts), + {ok, IDs} ?= query_recipients(ProcID, From, To, Opts), + {ok, base_layer_offsets(IDs, Opts)} + end. + +%% @doc Ensure the node's local Arweave index covers the block range, so that +%% every transaction in it has both a weave offset and a locally-cached header. +%% `~copycat@1.0/arweave' in `shallow' mode records an ID->offset entry for each +%% transaction and caches its (data-free) header -- populating the `field-target' +%% match index that `query_recipients' reads. `reindex=false' skips blocks the +%% node has already indexed (for any process), so the run is idempotent, +%% incremental, and shared: overlapping ranges across processes are not +%% re-fetched. +ensure_offsets(From, To, _Opts) when From > To -> ok; +ensure_offsets(From, To, Opts) -> + maybe + {ok, _} ?= + hb_ao:resolve( + << + "~copycat@1.0/arweave&mode=shallow&reindex=false", + "&from=", (hb_util:bin(To))/binary, + "&to=", (hb_util:bin(From))/binary + >>, + no_result_cache(Opts) + ), + ok + end. + +%% @doc Query for the transactions addressed to a process within a block +%% range, ordered by weave position, returning their IDs. By default the query +%% is served by the node's own `~query@1.0' index (`local'), whose `field-target' +%% matches are populated by `~copycat@1.0/arweave' as it indexes the range (see +%% `ensure_offsets'); a node may instead be configured +%% (`arweave_scheduler_query_source => remote') to query a remote gateway. +query_recipients(ProcID, From, To, Opts) -> + Source = hb_opts:get(arweave_scheduler_query_source, local, Opts), + Query = + << + "query($after: String) { transactions(", + "recipients: [\"", (hb_util:human_id(ProcID))/binary, "\"], ", + "block: { min: ", (hb_util:bin(From))/binary, + ", max: ", (hb_util:bin(To))/binary, " }, ", + "sort: HEIGHT_ASC, ", + "first: ", (hb_util:bin(?QUERY_PAGE_SIZE))/binary, + ", after: $after", + ") { pageInfo { hasNextPage } edges { cursor node { id } } } }" + >>, + query_pages(Source, Query, undefined, [], Opts). + +query_pages(Source, Query, After, Acc, Opts) -> + Variables = + case After of + undefined -> #{}; + _ -> #{ <<"after">> => After } + end, + maybe + {ok, Transactions} ?= run_query(Source, Query, Variables, Opts), + Edges = hb_maps:get(<<"edges">>, Transactions, [], Opts), + IDs = Acc ++ [ edge_id(E, Opts) || E <- Edges ], + HasNext = + hb_util:atom( + hb_ao:get( + <<"pageInfo/hasNextPage">>, + Transactions, + false, + Opts#{ <<"hashpath">> => ignore } + ) + ), + case {HasNext, Edges} of + {true, [_ | _]} -> + Cursor = + hb_maps:get(<<"cursor">>, lists:last(Edges), undefined, Opts), + query_pages(Source, Query, Cursor, IDs, Opts); + _ -> + {ok, [ ID || ID <- IDs, is_binary(ID) ]} + end + end. + +edge_id(Edge, Opts) -> + hb_maps:get(<<"id">>, hb_maps:get(<<"node">>, Edge, #{}, Opts), undefined, Opts). + +%% @doc Run a GraphQL `transactions' query, returning its connection. `local' +%% resolves the node's own `~query@1.0/graphql' endpoint; `remote' posts to +%% the configured `gateway' directly. The remote path deliberately does not use +%% `hb_client_gateway:query': its admissibility gate rejects legitimately-empty +%% ranges, and the racing `/graphql' route lets AO-search gateways that do not +%% index arbitrary L1 transactions win with empty-but-200 responses. Here an +%% empty range is an authoritative answer. +run_query(local, Query, Variables, Opts) -> + to_transactions( + hb_ao:resolve( + #{ <<"device">> => <<"query@1.0">> }, + #{ + <<"path">> => <<"graphql">>, + <<"method">> => <<"POST">>, + <<"query">> => Query, + <<"variables">> => Variables + }, + no_result_cache(Opts) + ), + Opts + ); +run_query(remote, Query, Variables, Opts) -> + Gateway = hb_opts:get(gateway, <<"https://arweave.net">>, Opts), + to_transactions( + hb_http:post( + Gateway, + #{ + <<"path">> => <<"/graphql">>, + <<"content-type">> => <<"application/json">>, + <<"body">> => + hb_json:encode( + #{ <<"query">> => Query, <<"variables">> => Variables } + ) + }, + no_result_cache(Opts) + ), + Opts + ). + +to_transactions({ok, Response}, Opts) -> + Decoded = + hb_json:decode(hb_ao:get(<<"body">>, Response, <<"{}">>, Opts)), + {ok, + hb_ao:get( + <<"data/transactions">>, + Decoded, + #{}, + Opts#{ <<"hashpath">> => ignore } + ) + }; +to_transactions({error, Reason}, _Opts) -> + {error, + #{ + <<"status">> => 502, + <<"reason">> => <<"Arweave query failed.">>, + <<"detail">> => Reason + } + }. + +%% @doc Annotate each matched transaction with its weave offset from the +%% local index, keeping only the base-layer (`tx@1.0') transactions and +%% dropping bundled data items (indexed under the `ans104@1.0' codec) and any +%% that are not yet locally indexed. The result is sorted by offset, which is +%% the canonical weave order. +base_layer_offsets(IDs, Opts) -> + Store = hb_store_arweave:store_from_opts(Opts), + lists:keysort( + 1, + lists:filtermap( + fun(ID) -> base_layer_offset(Store, ID, Opts) end, + IDs + ) + ). + +base_layer_offset(Store, ID, Opts) -> + case offset_entry(Store, ID, Opts) of + {ok, <<"tx@1.0">>, Offset} -> {true, {Offset, ID}}; + _ -> false + end. + +%% @doc Read a transaction's local index entry, returning its codec device and +%% weave offset. +offset_entry(Store, ID, Opts) -> + case hb_store_arweave:read_offset(Store, ID, Opts) of + {ok, + #{ + <<"codec-device">> := Codec, + <<"start-offset">> := Offset + }} when is_integer(Offset) -> + {ok, Codec, Offset}; + _ -> not_found + end. + +%% @doc Materialize a chunk's assignments and record that the range is synced +%% to `To'. The discovered list holds only the messages in blocks above the +%% already-synced range, so they are exactly the next slots -- appended from +%% `next-slot' with no de-duplication. `next-slot' and `synced-to' are then +%% persisted together, so the two never diverge. +materialize(ProcID, State = #{ <<"next-slot">> := NextSlot }, Ordered, To, Opts) -> + assign(ProcID, State, NextSlot, Ordered, To, Opts). + +assign(ProcID, State, Slot, [], To, Opts) -> + NewState = State#{ <<"next-slot">> => Slot, <<"synced-to">> => To }, + ok = dev_arweave_scheduler_cache:write_state(ProcID, NewState, Opts), + {ok, NewState}; +assign(ProcID, State, Slot, [{Offset, TXID} | Rest], To, Opts) -> + case read_tx_header(TXID, Opts) of + {ok, Msg} -> + ok = write_assignment(ProcID, Slot, Offset, Msg, Opts), + assign(ProcID, State, Slot + 1, Rest, To, Opts); + {error, Err} -> + ?event( + error, + {arweave_scheduler_tx_unavailable, + {proc_id, ProcID}, + {tx, {explicit, TXID}} + } + ), + {error, Err} + end. + +%% @doc Generate and store the synthetic assignment for a message. Mirrors the +%% assignments minted by `~scheduler@1.0', but the on-chain position is the +%% transaction's weave `offset' rather than a scheduler-assigned nonce. Every +%% field derives from chain data, so the assignment is deterministic across +%% nodes and is left uncommitted. +write_assignment(ProcID, Slot, Offset, Msg, Opts) -> + BaseAssignment = + lib_scheduler:base_assignment( + hb_util:human_id(ProcID), + Slot, + Msg, + Opts + ), + Assignment = BaseAssignment#{ <<"offset">> => Offset }, + ?event( + {minting_assignment, + {proc_id, ProcID}, + {slot, Slot}, + {offset, Offset} + } + ), + dev_arweave_scheduler_cache:write(Assignment, Opts). + +%% @doc Find or create the persisted synchronization state for a process. +ensure_initialized(ProcID, Opts) -> + case dev_arweave_scheduler_cache:read_state(ProcID, Opts) of + {ok, State} -> {ok, State}; + _ -> initialize(ProcID, Opts) + end. + +%% @doc First contact with a process: locate its spawn block, index it, read +%% the process header, and mint the slot 0 assignment from the process +%% message itself. The canonical process header is also written to the cache +%% so that `~process@1.0' resolves to the verifying tx@1.0 form +%% ahead of any lossier gateway-derived copy. If the spawn is not yet +%% confirmed, initialization fails and is retried on the next synchronization. +initialize(ProcID, Opts) -> + ?event({initializing_arweave_schedule, {proc_id, ProcID}}), + maybe + {ok, SpawnHeight} ?= spawn_height(ProcID, Opts), + ok ?= ensure_offsets(SpawnHeight, SpawnHeight, Opts), + {ok, Offset} ?= tx_offset(ProcID, Opts), + {ok, Process} ?= read_tx_header(ProcID, Opts), + {ok, _} = hb_cache:write(Process, Opts), + ok = write_assignment(ProcID, 0, Offset, Process, Opts), + % `synced-to' starts one below the spawn block: slot 0 is the process + % itself, and no message-bearing block has been indexed yet. The first + % sync begins its range at the spawn block, catching any messages mined + % alongside the process. + State = + #{ + <<"next-slot">> => 1, + <<"spawn-height">> => SpawnHeight, + <<"synced-to">> => SpawnHeight - 1 + }, + ok = dev_arweave_scheduler_cache:write_state(ProcID, State, Opts), + {ok, State} + end. + +%% @doc Read an L1 transaction as a header-only message from the node's stores. +%% `~copycat@1.0/arweave' caches the (data-free) header locally while indexing, +%% so the read is normally served straight from the local store; a miss falls +%% through the rest of the store chain (which may reach the header faster than a +%% specific Arweave host). The header carries the transaction's tags and its +%% tx@1.0 signature -- so its committed ID is the transaction ID and it verifies +%% -- and no data is attached at this layer, so the schedule never depends on +%% the availability of any transaction's data. +read_tx_header(TXID, Opts) -> + case hb_cache:read(TXID, Opts) of + {ok, Msg} -> {ok, hb_cache:ensure_all_loaded(Msg, Opts)}; + _ -> + {error, + #{ + <<"status">> => 503, + <<"reason">> => + <<"Transaction header is not yet retrievable ", + "from Arweave.">>, + <<"tx">> => TXID + } + } + end. + +%% @doc Read a transaction's weave offset from the local index. A transaction +%% that has not yet been indexed has no offset. +tx_offset(TXID, Opts) -> + Store = hb_store_arweave:store_from_opts(Opts), + case offset_entry(Store, TXID, Opts) of + {ok, _Codec, Offset} -> {ok, Offset}; + not_found -> + {error, + #{ + <<"status">> => 503, + <<"reason">> => + <<"Transaction is not yet indexed locally.">>, + <<"tx">> => TXID + } + } + end. + +%% @doc Locate the block height that includes a transaction, via the Arweave +%% `tx status' API. An unconfirmed transaction has no block height: a process +%% cannot be scheduled against until its spawn is confirmed. +spawn_height(TXID, Opts) -> + Res = + hb_http:request( + #{ + <<"method">> => <<"GET">>, + <<"path">> => <<"/arweave/tx/", TXID/binary, "/status">> + }, + no_result_cache(Opts) + ), + case lib_arweave_common:best_response(Res) of + {ok, Response} -> + Status = + hb_json:decode(hb_ao:get(<<"body">>, Response, <<"{}">>, Opts)), + case hb_maps:get(<<"block_height">>, Status, undefined, Opts) of + undefined -> {error, unconfirmed_process(TXID)}; + Height -> {ok, hb_util:int(Height)} + end; + _ -> {error, unconfirmed_process(TXID)} + end. + +unconfirmed_process(TXID) -> + #{ + <<"status">> => 404, + <<"reason">> => <<"Process is not yet confirmed on Arweave.">>, + <<"process">> => TXID + }. + +%% @doc The height up to which the schedule may safely be extended: the +%% network tip less the configured confirmation depth, optionally capped by +%% the `arweave_scheduler_max_height' node option. +confirmed_tip(Opts) -> + maybe + {ok, Height} ?= + hb_ao:resolve( + <>, + no_result_cache(Opts) + ), + Depth = + hb_util:int( + hb_opts:get( + arweave_scheduler_confirmation_depth, + ?DEFAULT_CONFIRMATION_DEPTH, + Opts + ) + ), + Confirmed = hb_util:int(Height) - Depth, + case hb_opts:get(arweave_scheduler_max_height, undefined, Opts) of + undefined -> {ok, Confirmed}; + Max -> {ok, min(Confirmed, hb_util:int(Max))} + end + end. + +%%% Tests + +%%% The permanent test fixture: a `lua@5.3a' process scheduled by this +%%% device, seeded onto Arweave with three state-transformation messages +%%% (`test/arweave-scheduler-test.lua'; all transactions carry +%%% `test-suite: arweave-scheduler' tags). `?FIXTURE_MSG2' was posted +%%% directly to arweave.net, never passing through a HyperBEAM node: its +%%% presence in the schedule proves foreign-message indexing. Arweave is +%%% permanent, so these tests are repeatable against the live network. +%%% Spawn block 1958986; messages at 1958993, 1958994 and 1958995. +-define(FIXTURE_MODULE, <<"_GeSyZbQkmqWk6YzL-tjIqJ-2hakIkg-9k127DpVfO8">>). +-define(FIXTURE_PROCESS, <<"q3SycbYpO1lz-S6V2kd7FG3DIZ3AU0NrKKxWw4C3yos">>). +-define(FIXTURE_MSG1, <<"Y8GnKytC57VUk7W7FlGnD1oH4OAwFHyP-pJPGCJBa3Y">>). +-define(FIXTURE_MSG2, <<"luf1fFmhi0RMIZNnFmv1QgK2SJU5plAFQ3ZxndUsLpk">>). +-define(FIXTURE_MSG3, <<"UThcMtT6zy0pJ7iXUsMTFUMu3Xhv51EAjn5BzqHMhR4">>). +-define(FIXTURE_MAX_HEIGHT, 1958995). + +test_opts() -> + TestStore = hb_test_utils:test_store(), + IndexStore = hb_test_utils:test_store(), + % The full default config is merged in so the node has the Arweave routes + % that `~copycat@1.0' and `~arweave@2.9' resolve through. + (hb_opts:default_message())#{ + <<"store">> => [ + TestStore, + % The local Arweave offset index that discovery orders by. + #{ + <<"store-module">> => hb_store_arweave, + <<"name">> => <<"cache-arweave">>, + <<"index-store">> => [IndexStore] + }, + % Serves the Lua module transaction's data (its source) and the + % transaction headers `~copycat@1.0/arweave' fetches while indexing. + #{ + <<"store-module">> => hb_store_gateway, + <<"local-store">> => [TestStore] + } + ], + <<"arweave-index-store">> => #{ <<"index-store">> => [IndexStore] }, + <<"arweave-index-workers">> => 8, + <<"arweave-scheduler-confirmation-depth">> => 1, + <<"arweave-scheduler-max-height">> => ?FIXTURE_MAX_HEIGHT, + <<"priv-wallet">> => ar_wallet:new() + }. + +%% @doc Scheduling a message that does not carry a `tx@1.0' commitment must +%% be rejected. +reject_non_tx_message_test() -> + Opts = #{ <<"priv-wallet">> => ar_wallet:new() }, + Msg = + hb_message:commit( + #{ + <<"target">> => hb_util:human_id(crypto:strong_rand_bytes(32)), + <<"test-key">> => <<"test-value">> + }, + Opts + ), + ?assertMatch( + {error, #{ <<"status">> := 422 }}, + hb_ao:resolve( + #{ <<"device">> => <<"arweave-scheduler@1.0">> }, + #{ + <<"path">> => <<"schedule">>, + <<"method">> => <<"POST">>, + <<"body">> => Msg + }, + Opts + ) + ). + +%% @doc A validly `tx@1.0'-committed message that carries data must also be +%% rejected: this scheduler sequences transaction headers only, so it will not +%% dispatch (nor later depend on the chunk availability of) a message's data. +reject_data_message_test() -> + Opts = #{ <<"priv-wallet">> => ar_wallet:new() }, + Msg = + hb_message:commit( + #{ + <<"target">> => hb_util:human_id(crypto:strong_rand_bytes(32)), + <<"data">> => <<"unschedulable-data">> + }, + Opts, + #{ <<"commitment-device">> => <<"tx@1.0">> } + ), + ?assertMatch( + {error, #{ <<"status">> := 422 }}, + hb_ao:resolve( + #{ <<"device">> => <<"arweave-scheduler@1.0">> }, + #{ + <<"path">> => <<"schedule">>, + <<"method">> => <<"POST">>, + <<"body">> => Msg + }, + Opts + ) + ). + +%% @doc Base-layer annotation keeps only tx@1.0 offsets and sorts by offset. +base_layer_offsets_test() -> + Store = hb_test_utils:test_store(), + hb_store:start(Store), + ArwStore = #{ <<"index-store">> => [Store] }, + Opts = #{ <<"arweave-index-store">> => ArwStore }, + Late = hb_util:human_id(crypto:strong_rand_bytes(32)), + Early = hb_util:human_id(crypto:strong_rand_bytes(32)), + Bundled = hb_util:human_id(crypto:strong_rand_bytes(32)), + Unindexed = hb_util:human_id(crypto:strong_rand_bytes(32)), + ok = hb_store_arweave:write_offset(ArwStore, Late, <<"tx@1.0">>, 200, 1), + ok = hb_store_arweave:write_offset(ArwStore, Early, <<"tx@1.0">>, 100, 1), + ok = hb_store_arweave:write_offset(ArwStore, Bundled, <<"ans104@1.0">>, 150, 1), + ?assertEqual( + [{100, Early}, {200, Late}], + base_layer_offsets([Late, Bundled, Early, Unindexed], Opts) + ). + +%% @doc `/status' reports each tracked process's contiguously-indexed block +%% range straight from the cache, without triggering a synchronization. +status_attestable_range_test() -> + Store = hb_test_utils:test_store(), + hb_store:start(Store), + Opts = #{ <<"store">> => [Store], <<"priv-wallet">> => ar_wallet:new() }, + ProcID = hb_util:human_id(crypto:strong_rand_bytes(32)), + ok = + dev_arweave_scheduler_cache:write_state( + ProcID, + #{ + <<"next-slot">> => 4, + <<"spawn-height">> => 100, + <<"synced-to">> => 150 + }, + Opts + ), + {ok, #{ <<"processes">> := Processes }} = status(#{}, #{}, Opts), + ?assertEqual( + #{ <<"from">> => 100, <<"to">> => 150, <<"current">> => 3 }, + hb_maps:get(ProcID, Processes, not_found, Opts) + ). + +%% @doc Synchronize the fixture schedule from the live network, retrying on +%% transient indexing failures (for example, gateway rate limits while +%% fetching block transaction headers). +fixture_sync(_Opts, 0) -> {error, fixture_sync_failed}; +fixture_sync(Opts, Attempts) -> + Res = + hb_ao:resolve( + #{ <<"device">> => <<"arweave-scheduler@1.0">> }, + #{ + <<"path">> => <<"schedule">>, + <<"method">> => <<"GET">>, + <<"target">> => ?FIXTURE_PROCESS + }, + Opts + ), + case Res of + {ok, Schedule} -> {ok, Schedule}; + {error, _} -> + timer:sleep(5000), + fixture_sync(Opts, Attempts - 1) + end. + +%% @doc Read the fixture schedule from the live network and check that the +%% assignments arrive in seeded order, including the message that never +%% passed through a HyperBEAM node. +fixture_schedule_test_() -> + {timeout, 1200, fun fixture_schedule/0}. +fixture_schedule() -> + Opts = test_opts(), + {ok, Schedule} = fixture_sync(Opts, 5), + Assignments = + hb_ao:normalize_keys( + hb_ao:get(<<"assignments">>, Schedule, Opts), + Opts + ), + SlotIDs = + lists:map( + fun(Slot) -> + Assignment = + hb_maps:get(hb_util:bin(Slot), Assignments, not_found, Opts), + ?assertEqual( + Slot, + hb_util:int(hb_ao:get(<<"slot">>, Assignment, Opts)) + ), + Body = hb_ao:get(<<"body">>, Assignment, Opts), + % The body is a header: it verifies and its committed ID is + % the transaction ID, but it carries no data payload -- the + % schedule never depends on data availability. + ?assert(hb_message:verify(Body, all, Opts)), + BodyKeys = hb_maps:keys(hb_message:uncommitted(Body, Opts), Opts), + ?assertNot(lists:member(<<"data">>, BodyKeys)), + ?assertNot(lists:member(<<"body">>, BodyKeys)), + hb_message:id(Body, signed, Opts) + end, + [0, 1, 2, 3] + ), + ?assertEqual( + [?FIXTURE_PROCESS, ?FIXTURE_MSG1, ?FIXTURE_MSG2, ?FIXTURE_MSG3], + SlotIDs + ). + +%% @doc Compute the fixture process from its Arweave schedule and check the +%% state transformations applied in order: `setstate 1000', `addstate 337', +%% then `querystate' reporting the result. +fixture_lua_e2e_test_() -> + {timeout, 1200, fun fixture_lua_e2e/0}. +fixture_lua_e2e() -> + Opts = test_opts(), + % Prime the schedule first: synchronization indexes the spawn block, so + % the process message is read back as its canonical tx@1.0 decoding + % (rather than the gateway store's lossy representation). + {ok, _} = fixture_sync(Opts, 5), + {ok, RawProcess} = hb_cache:read(?FIXTURE_PROCESS, Opts), + Process = hb_cache:ensure_all_loaded(RawProcess, Opts), + % Luerl represents Lua numbers as floats. + ?assertEqual( + {ok, 1000.0}, + hb_ao:resolve( + Process, + #{ <<"path">> => <<"compute/state">>, <<"slot">> => 1 }, + Opts + ) + ), + ?assertEqual( + {ok, 1337.0}, + hb_ao:resolve( + Process, + #{ <<"path">> => <<"compute/state">>, <<"slot">> => 2 }, + Opts + ) + ), + ?assertEqual( + {ok, <<"state=1337.0">>}, + hb_ao:resolve( + Process, + #{ + <<"path">> => <<"compute/results/output/body">>, + <<"slot">> => 3 + }, + Opts + ) + ), + ?assertEqual({ok, 1337.0}, hb_ao:resolve(Process, <<"now/state">>, Opts)). + +%% @doc A first sync indexes the range in bounded chunks, persisting its +%% progress after each. With a chunk size smaller than the fixture's span -- +%% so the messages fall in different chunks -- the schedule is still assembled +%% in order, and `/status' reports the range synced all the way to the tip. +chunked_sync_test_() -> + {timeout, 1200, fun chunked_sync/0}. +chunked_sync() -> + Opts = (test_opts())#{ <<"arweave-scheduler-sync-chunk">> => 4 }, + {ok, Schedule} = fixture_sync(Opts, 5), + Assignments = + hb_ao:normalize_keys( + hb_ao:get(<<"assignments">>, Schedule, Opts), + Opts + ), + SlotIDs = + lists:map( + fun(Slot) -> + Assignment = + hb_maps:get(hb_util:bin(Slot), Assignments, not_found, Opts), + hb_message:id(hb_ao:get(<<"body">>, Assignment, Opts), signed, Opts) + end, + [0, 1, 2, 3] + ), + ?assertEqual( + [?FIXTURE_PROCESS, ?FIXTURE_MSG1, ?FIXTURE_MSG2, ?FIXTURE_MSG3], + SlotIDs + ), + % A completed sync makes the whole fixture range attestable: the spawn + % block through the tip. + {ok, #{ <<"processes">> := Processes }} = status(#{}, #{}, Opts), + Range = hb_maps:get(?FIXTURE_PROCESS, Processes, not_found, Opts), + ?assertEqual(?FIXTURE_MAX_HEIGHT, hb_util:int(hb_ao:get(<<"to">>, Range, Opts))), + ?assertEqual(3, hb_util:int(hb_ao:get(<<"current">>, Range, Opts))). diff --git a/src/preloaded/process/dev_arweave_scheduler_cache.erl b/src/preloaded/process/dev_arweave_scheduler_cache.erl new file mode 100644 index 000000000..a8af37424 --- /dev/null +++ b/src/preloaded/process/dev_arweave_scheduler_cache.erl @@ -0,0 +1,111 @@ +%%% @doc A cache for the `~arweave-scheduler@1.0' device: stores the +%%% synthetic assignments derived from Arweave L1 transactions, as well as +%%% the per-process synchronization state of the indexer. +-module(dev_arweave_scheduler_cache). +-export([write/2, read/3, list_processes/1]). +-export([read_state/2, write_state/3]). +-export([assignments_to_bundle/4]). +-include("include/hb.hrl"). +-include_lib("eunit/include/eunit.hrl"). + +%%% The pseudo-path prefix which the arweave-scheduler cache should use. +-define(CACHE_PREFIX, <<"~arweave-scheduler@1.0">>). + +%% @doc Merge the scheduler store with the main store. Used before writing +%% to the cache. +opts(Opts) -> lib_scheduler:cache_opts(Opts). + +%% @doc Write an assignment message into the cache. Assignments are +%% deterministic derivations of chain data, so writes are idempotent: +%% concurrent synchronizations of the same process converge on identical +%% messages at identical paths. +write(Assignment, Opts) -> + lib_scheduler:write_assignment(?CACHE_PREFIX, Assignment, Opts). + +%% @doc Get an assignment message from the cache. +read(ProcID, Slot, Opts) -> + lib_scheduler:read_assignment(?CACHE_PREFIX, ProcID, Slot, Opts). + +%% @doc List the processes that the device holds synchronization state for. +list_processes(RawOpts) -> + Opts = opts(RawOpts), + Store = hb_opts:get(store, no_viable_store, Opts), + case hb_store:list(Store, <>, Opts) of + {ok, Names} -> Names; + _ -> [] + end. + +%% @doc Read the persisted synchronization state for a process. Returns +%% `{ok, State}' with the `next-slot' to be assigned, the process's +%% `spawn-height', and `synced-to' (the highest block whose messages have been +%% contiguously indexed and materialized), or propagates the store's +%% `not_found'. +read_state(ProcID, RawOpts) -> + Opts = opts(RawOpts), + Store = hb_opts:get(store, no_viable_store, Opts), + maybe + {ok, NextSlot} ?= + hb_store:read(Store, state_path(ProcID, <<"next-slot">>), Opts), + {ok, SpawnHeight} ?= + hb_store:read(Store, state_path(ProcID, <<"spawn-height">>), Opts), + {ok, SyncedTo} ?= + hb_store:read(Store, state_path(ProcID, <<"synced-to">>), Opts), + {ok, + #{ + <<"next-slot">> => hb_util:int(NextSlot), + <<"spawn-height">> => hb_util:int(SpawnHeight), + <<"synced-to">> => hb_util:int(SyncedTo) + } + } + end. + +%% @doc Persist the synchronization state for a process. +write_state(ProcID, State, RawOpts) -> + Opts = opts(RawOpts), + Store = hb_opts:get(store, no_viable_store, Opts), + #{ + <<"next-slot">> := NextSlot, + <<"spawn-height">> := SpawnHeight, + <<"synced-to">> := SyncedTo + } = State, + hb_store:write( + Store, + #{ + state_path(ProcID, <<"next-slot">>) => hb_util:bin(NextSlot), + state_path(ProcID, <<"spawn-height">>) => hb_util:bin(SpawnHeight), + state_path(ProcID, <<"synced-to">>) => hb_util:bin(SyncedTo) + }, + Opts + ). + +state_path(ProcID, Key) -> + hb_path:to_binary( + [ + ?CACHE_PREFIX, + <<"state">>, + hb_util:human_id(ProcID), + Key + ] + ). + +%% @doc Generate a `GET /schedule' response for a process. Mirrors +%% `dev_scheduler_formats:assignments_to_bundle/4' (which cannot be called +%% across device namespaces), less the bundle-level `timestamp', +%% `block-height' and `block-hash'. Those describe the current weave tip +%% rather than the blocks that sequenced these assignments, so they would be +%% misleading on a schedule that is a deterministic read of historical chain +%% data. The on-chain position of each message is its assignment's `offset'. +assignments_to_bundle(ProcID, Assignments, More, RawOpts) -> + Opts = lib_scheduler:format_opts(RawOpts), + {ok, #{ + <<"type">> => <<"schedule">>, + <<"process">> => hb_util:human_id(ProcID), + <<"continues">> => hb_util:atom(More), + <<"assignments">> => + hb_message:normalize_commitments( + hb_maps:from_list( + [ {hb_ao:get(<<"slot">>, A, Opts), A} || A <- Assignments ] + ), + Opts + ) + }}. diff --git a/test/arweave-scheduler-test.lua b/test/arweave-scheduler-test.lua new file mode 100644 index 000000000..ec816ccb4 --- /dev/null +++ b/test/arweave-scheduler-test.lua @@ -0,0 +1,34 @@ +--- @module arweave-scheduler-test +--- State-transformation test module for the `~arweave-scheduler@1.0` device. +--- Messages scheduled as Arweave L1 transactions invoke these functions via +--- their `path` tag. The function names deliberately avoid the keys that +--- `~lua@5.3a` excludes from its default handler (`set`, `path`, ...) and +--- the keys present on the process message itself. + +--- Executed for assignments without a `path` tag -- most notably slot 0, +--- the process message itself. +function compute(base, req, opts) + return base +end + +--- Overwrite the process state with the `value` tag of the message. +function setstate(base, req, opts) + base.state = tonumber(req.body.value) + return base +end + +--- Add the `value` tag of the message to the process state. +function addstate(base, req, opts) + base.state = (base.state or 0) + tonumber(req.body.value) + return base +end + +--- Report the process state in the results of the slot. +function querystate(base, req, opts) + base.results = { + output = { + body = "state=" .. tostring(base.state) + } + } + return base +end From 1fd8d8551dc1cff1ce4e3ac06ba4e5ce5970ffec Mon Sep 17 00:00:00 2001 From: Sam Williams Date: Fri, 24 Jul 2026 23:00:09 -0400 Subject: [PATCH 6/7] feat(arweave-scheduler): sequence a process by every L1 transaction A process message may now carry `scheduler-mode: all', and its schedule becomes every base-layer transaction on the network rather than only those whose `target' is the process. This is what lets a process observe value moving between two other addresses -- a payment it is owed but is not a party to -- which no recipient-filtered schedule can see. The mode is read from the process message itself at first contact and pinned into the persisted sync state. The process message is an L1 transaction, so the mode is chain data like the rest of the schedule, and every node derives the same one; a node option would have made slot -> body a function of node config, producing two indistinguishable, equally `valid' schedules for one process id. An unrecognized mode falls back to `target' rather than erroring: a spawn tag cannot be corrected once it is on-chain. `all' mode does not query. Dropping `recipients:' from the GraphQL would have silently returned nothing -- `dev_query_arweave' only builds candidate sets from height/id/ids/tags/owners/recipients, and `block:' is a post-filter over a set that must already exist -- so the range's block headers are walked directly instead, taking each block's own `txs' list. That is the same source `~copycat@1.0/arweave' enumerates, and `ensure_offsets' has already cached those blocks, so it is a local read. Nor does it re-sort by weave offset. Offset is not a total order: `TXStartOffset = TXEndOffset - data_size' and the accumulator only advances by `data_size', so every data-free transaction in a block shares the offset of the one before it -- and a data-free header is exactly what an AO message is. Block order (heights ascending, then each block's own transaction order) is a total order and is already canonical. The process's own transaction is skipped, as the first sync deliberately re-scans the spawn block and slot 0 is already the process. Assignments in this mode also record the `block-height' that sequenced them, so a process reading the whole chain has a deterministic clock; until now an assignment carried only its offset, and a block-denominated deadline was inexpressible. `target'-mode assignments are untouched, so existing schedules keep their assignment ids. Two things this mode cannot do the way `target' mode does. It must not drop a transaction the block lists but the local index does not hold: indexing a block is all-or-nothing and reports success either way, so dropping them would silently shorten the schedule and -- slots being positional -- shift every later slot on that node alone, with both nodes serving 200s and claiming the same attestable range. The range fails instead, leaving `synced-to' for the next pass to retry. And it sequences bundle transactions, whose headers the shallow indexing pass deliberately does not cache, so a header that is not held locally is fetched data-free on demand rather than stopping the sync. Tested: mode selection and its fallback, slot-0 detail per mode, block enumeration order and process-tx exclusion, offset annotation without re-sorting, failing closed on an unindexed transaction, the assignment cache round trip carrying block-height, and the mode's round trip through the sync state. The existing live-weave fixture tests still pass unchanged. --- .../process/dev_arweave_scheduler.erl | 363 ++++++++++++++++-- .../process/dev_arweave_scheduler_cache.erl | 19 +- 2 files changed, 354 insertions(+), 28 deletions(-) diff --git a/src/preloaded/process/dev_arweave_scheduler.erl b/src/preloaded/process/dev_arweave_scheduler.erl index 37a32ea45..3977ef3a0 100644 --- a/src/preloaded/process/dev_arweave_scheduler.erl +++ b/src/preloaded/process/dev_arweave_scheduler.erl @@ -27,6 +27,26 @@ %%% the device's own schedule cache. %%% %%% +%%% A process message may widen what it is sequenced by, with a +%%% `scheduler-mode' key: +%%%
    +%%%
  • `target' (the default): the process's messages are the transactions +%%% addressed to it, as above.
  • +%%%
  • `all': every base-layer transaction is a message in the +%%% process's schedule, whoever it is addressed to. This is what lets a +%%% process observe value moving between two other addresses -- a +%%% payment it is owed but is not a party to -- at the cost of a slot +%%% for every transaction on the network. Its schedule is enumerated +%%% from the block headers rather than by query, in canonical chain +%%% order (blocks ascending, then each block's own transaction order), +%%% and each assignment records the `block-height' that sequenced it, so +%%% a process in this mode has a clock. The process's own transaction is +%%% not re-assigned: it remains slot 0 alone.
  • +%%%
+%%% The mode is read from the process message itself -- an L1 transaction -- +%%% so it is chain data like the rest of the schedule, and is fixed for the +%%% life of the process. +%%% %%% Assignment bodies are transaction headers only: each is the %%% data-free header that `~copycat@1.0/arweave' cached while indexing the %%% range (read from the node's stores), so the schedule never carries (nor @@ -93,6 +113,9 @@ -define(QUERY_PAGE_SIZE, 100). %%% The Arweave device that we resolve chain data through. -define(ARWEAVE_DEVICE, <<"~arweave@2.9">>). +%%% The sequencing mode of a process whose message does not name one: only the +%%% transactions addressed to it are its messages. +-define(DEFAULT_MODE, <<"target">>). %% @doc This device uses a default handler to route requests to the correct %% function. @@ -364,19 +387,21 @@ sync(ProcID, Opts) -> do_sync(_ProcID, State = #{ <<"synced-to">> := SyncedTo }, Upper, _Opts) when SyncedTo >= Upper -> {ok, State}; -do_sync(ProcID, State = #{ <<"synced-to">> := SyncedTo }, Upper, Opts) -> +do_sync(ProcID, State = #{ <<"synced-to">> := SyncedTo, <<"mode">> := Mode }, + Upper, Opts) -> From = SyncedTo + 1, To = min(SyncedTo + sync_chunk_blocks(Opts), Upper), ?event( {arweave_scheduler_sync, {proc_id, ProcID}, + {mode, Mode}, {from, From}, {to, To}, {target, Upper} } ), maybe - {ok, Ordered} ?= discover(ProcID, From, To, Opts), + {ok, Ordered} ?= discover(ProcID, Mode, From, To, Opts), {ok, NewState} ?= materialize(ProcID, State, Ordered, To, Opts), do_sync(ProcID, NewState, Upper, Opts) end. @@ -397,21 +422,68 @@ no_result_cache(Opts) -> <<"cache-control">> => [<<"no-cache">>, <<"no-store">>] }. -%% @doc Discover the base-layer transactions addressed to a process within a -%% block range, in ascending weave order, entirely from the node's own index. -%% Indexing the range with `~copycat@1.0/arweave' both records each +%% @doc Discover the base-layer transactions that a process is sequenced by +%% within a block range, in canonical weave order, entirely from the node's own +%% index. Indexing the range with `~copycat@1.0/arweave' both records each %% transaction's weave offset and caches its header (so its `target' is locally -%% matchable); the recipient match is then served from the node's own +%% matchable), whichever mode the process is in. +%% +%% In `target' mode the recipient match is served from the node's own %% `~query@1.0' GraphQL endpoint, and each match is annotated with its offset -- %% the sort key and the `offset' recorded on the assignment. No gateway is -%% queried by default. Returns `{Offset, TXID}' pairs in ascending offset order; -%% bundled data items are excluded, as this scheduler sequences the base layer -%% only. -discover(ProcID, From, To, Opts) -> +%% queried by default. In `all' mode there is nothing to match, so no query is +%% run: the range's block headers are walked directly (see +%% `enumerate_blocks/4'). +%% +%% Returns `{Extra, TXID}' pairs in the order they are to be assigned, where +%% `Extra' is the sequencing detail recorded on each assignment. Bundled data +%% items are excluded, as this scheduler sequences the base layer only. +discover(ProcID, <<"all">>, From, To, Opts) -> + maybe + ok ?= ensure_offsets(From, To, Opts), + {ok, Located} ?= enumerate_blocks(ProcID, From, To, Opts), + base_layer_blocks(Located, Opts) + end; +discover(ProcID, _Mode, From, To, Opts) -> maybe ok ?= ensure_offsets(From, To, Opts), {ok, IDs} ?= query_recipients(ProcID, From, To, Opts), - {ok, base_layer_offsets(IDs, Opts)} + {ok, + [ + {#{ <<"offset">> => Offset }, TXID} + || + {Offset, TXID} <- base_layer_offsets(IDs, Opts) + ] + } + end. + +%% @doc Enumerate every transaction in a block range from the block headers +%% themselves, in canonical chain order: blocks ascending by height, then each +%% block's transactions in the order that block lists them. This is the same +%% source `~copycat@1.0/arweave' walks while indexing, and `ensure_offsets' has +%% already cached the blocks locally, so the enumeration is normally a local +%% read. The process's own transaction is skipped: it is already slot 0. +%% +%% Returns `{Height, TXID}' pairs. Unlike a query result these are ordered by +%% construction, and that order is total -- weave offset is not, because every +%% transaction that carries no data shares the offset of the one before it. +enumerate_blocks(_ProcID, From, To, _Opts) when From > To -> {ok, []}; +enumerate_blocks(ProcID, From, To, Opts) -> + maybe + {ok, Block} ?= + hb_ao:resolve( + <>, + no_result_cache(Opts) + ), + {ok, Rest} ?= enumerate_blocks(ProcID, From + 1, To, Opts), + {ok, + [ + {From, TXID} + || + TXID <- hb_maps:get(<<"txs">>, Block, [], Opts), + TXID =/= hb_util:human_id(ProcID) + ] ++ Rest + } end. %% @doc Ensure the node's local Arweave index covers the block range, so that @@ -570,6 +642,60 @@ base_layer_offset(Store, ID, Opts) -> _ -> false end. +%% @doc Annotate each block-enumerated transaction with its weave offset, +%% keeping only the base-layer ones. The enumeration is already in canonical +%% order, so -- unlike `base_layer_offsets/2' -- it is not re-sorted. Each +%% assignment records the height of the block that sequenced it as well as its +%% offset: in this mode the schedule follows the chain rather than the process, +%% so the block height is the only clock a process has. +%% +%% A transaction the block lists but the local index does not hold is a failure +%% to index, not an absence -- `~copycat@1.0/arweave' skips a whole block if any +%% one of its transaction headers could not be fetched, while still reporting +%% success. Dropping those would silently shorten the schedule and, because +%% slots are positional, shift every later slot on this node alone. So the range +%% fails here instead, leaving `synced-to' where it was for the next pass to +%% retry. +base_layer_blocks(Located, Opts) -> + Store = hb_store_arweave:store_from_opts(Opts), + lists:foldr( + fun(_, {error, Err}) -> {error, Err}; + ({Height, ID}, {ok, Acc}) -> + case offset_entry(Store, ID, Opts) of + {ok, <<"tx@1.0">>, Offset} -> + {ok, + [ + { + #{ + <<"offset">> => Offset, + <<"block-height">> => Height + }, + ID + } + | + Acc + ] + }; + {ok, _Codec, _Offset} -> + % A bundled data item: this scheduler sequences the base + % layer only. + {ok, Acc}; + not_found -> + {error, + #{ + <<"status">> => 503, + <<"reason">> => + <<"Block range is not fully indexed locally.">>, + <<"tx">> => ID, + <<"block-height">> => Height + } + } + end + end, + {ok, []}, + Located + ). + %% @doc Read a transaction's local index entry, returning its codec device and %% weave offset. offset_entry(Store, ID, Opts) -> @@ -595,10 +721,10 @@ assign(ProcID, State, Slot, [], To, Opts) -> NewState = State#{ <<"next-slot">> => Slot, <<"synced-to">> => To }, ok = dev_arweave_scheduler_cache:write_state(ProcID, NewState, Opts), {ok, NewState}; -assign(ProcID, State, Slot, [{Offset, TXID} | Rest], To, Opts) -> +assign(ProcID, State, Slot, [{Extra, TXID} | Rest], To, Opts) -> case read_tx_header(TXID, Opts) of {ok, Msg} -> - ok = write_assignment(ProcID, Slot, Offset, Msg, Opts), + ok = write_assignment(ProcID, Slot, Extra, Msg, Opts), assign(ProcID, State, Slot + 1, Rest, To, Opts); {error, Err} -> ?event( @@ -613,10 +739,11 @@ assign(ProcID, State, Slot, [{Offset, TXID} | Rest], To, Opts) -> %% @doc Generate and store the synthetic assignment for a message. Mirrors the %% assignments minted by `~scheduler@1.0', but the on-chain position is the -%% transaction's weave `offset' rather than a scheduler-assigned nonce. Every -%% field derives from chain data, so the assignment is deterministic across -%% nodes and is left uncommitted. -write_assignment(ProcID, Slot, Offset, Msg, Opts) -> +%% transaction's weave `offset' rather than a scheduler-assigned nonce (joined, +%% in `all' mode, by the `block-height' that sequenced it). Every field derives +%% from chain data, so the assignment is deterministic across nodes and is left +%% uncommitted. +write_assignment(ProcID, Slot, Extra, Msg, Opts) -> BaseAssignment = lib_scheduler:base_assignment( hb_util:human_id(ProcID), @@ -624,12 +751,12 @@ write_assignment(ProcID, Slot, Offset, Msg, Opts) -> Msg, Opts ), - Assignment = BaseAssignment#{ <<"offset">> => Offset }, + Assignment = hb_maps:merge(BaseAssignment, Extra, Opts), ?event( {minting_assignment, {proc_id, ProcID}, {slot, Slot}, - {offset, Offset} + {extra, Extra} } ), dev_arweave_scheduler_cache:write(Assignment, Opts). @@ -655,7 +782,15 @@ initialize(ProcID, Opts) -> {ok, Offset} ?= tx_offset(ProcID, Opts), {ok, Process} ?= read_tx_header(ProcID, Opts), {ok, _} = hb_cache:write(Process, Opts), - ok = write_assignment(ProcID, 0, Offset, Process, Opts), + Mode = mode(Process, Opts), + ok = + write_assignment( + ProcID, + 0, + slot_zero(Mode, Offset, SpawnHeight), + Process, + Opts + ), % `synced-to' starts one below the spawn block: slot 0 is the process % itself, and no message-bearing block has been indexed yet. The first % sync begins its range at the spawn block, catching any messages mined @@ -664,12 +799,31 @@ initialize(ProcID, Opts) -> #{ <<"next-slot">> => 1, <<"spawn-height">> => SpawnHeight, - <<"synced-to">> => SpawnHeight - 1 + <<"synced-to">> => SpawnHeight - 1, + <<"mode">> => Mode }, ok = dev_arweave_scheduler_cache:write_state(ProcID, State, Opts), {ok, State} end. +%% @doc The sequencing mode a process message asks for. A mode the device does +%% not implement falls back to `target' rather than erroring: the process +%% message is immutable, so a typo in a spawn tag would otherwise wedge the +%% process permanently. +mode(Process, Opts) -> + case hb_ao:get(<<"scheduler-mode">>, Process, ?DEFAULT_MODE, Opts) of + <<"all">> -> <<"all">>; + _ -> ?DEFAULT_MODE + end. + +%% @doc The sequencing detail recorded on slot 0. The process message is its +%% own first message, so in `all' mode it carries the height of its spawn block +%% -- the process's clock starts at the block it was created in. +slot_zero(<<"all">>, Offset, SpawnHeight) -> + #{ <<"offset">> => Offset, <<"block-height">> => SpawnHeight }; +slot_zero(_Mode, Offset, _SpawnHeight) -> + #{ <<"offset">> => Offset }. + %% @doc Read an L1 transaction as a header-only message from the node's stores. %% `~copycat@1.0/arweave' caches the (data-free) header locally while indexing, %% so the read is normally served straight from the local store; a miss falls @@ -680,6 +834,27 @@ initialize(ProcID, Opts) -> %% the availability of any transaction's data. read_tx_header(TXID, Opts) -> case hb_cache:read(TXID, Opts) of + {ok, Msg} -> {ok, hb_cache:ensure_all_loaded(Msg, Opts)}; + _ -> fetch_tx_header(TXID, Opts) + end. + +%% @doc Fetch a transaction's data-free header from Arweave. The indexing pass +%% caches the headers of plain L1 transactions as it goes, but deliberately not +%% those of bundles -- and an `all'-mode schedule sequences every transaction in +%% a block, bundles included. Their headers are fetched on demand rather than +%% failing the whole range; `exclude-data' keeps the schedule header-only, so it +%% still never depends on the availability of a transaction's data. +fetch_tx_header(TXID, Opts) -> + Res = + hb_ao:resolve( + << + ?ARWEAVE_DEVICE/binary, "/tx", + "&tx=", TXID/binary, + "&exclude-data=true" + >>, + no_result_cache(Opts) + ), + case Res of {ok, Msg} -> {ok, hb_cache:ensure_all_loaded(Msg, Opts)}; _ -> {error, @@ -880,6 +1055,149 @@ base_layer_offsets_test() -> base_layer_offsets([Late, Bundled, Early, Unindexed], Opts) ). +%% @doc The sequencing mode is read from the process message, and anything the +%% device does not implement leaves the process sequenced as normal rather than +%% wedging it: the process message cannot be corrected once it is on-chain. +mode_test() -> + ?assertEqual(<<"all">>, mode(#{ <<"scheduler-mode">> => <<"all">> }, #{})), + ?assertEqual(?DEFAULT_MODE, mode(#{}, #{})), + ?assertEqual( + ?DEFAULT_MODE, + mode(#{ <<"scheduler-mode">> => <<"sideways">> }, #{}) + ). + +%% @doc A process sequenced by the whole chain has a clock from slot 0; one +%% sequenced by its own messages keeps the assignment shape it always had. +slot_zero_test() -> + ?assertEqual( + #{ <<"offset">> => 7, <<"block-height">> => 3 }, + slot_zero(<<"all">>, 7, 3) + ), + ?assertEqual(#{ <<"offset">> => 7 }, slot_zero(?DEFAULT_MODE, 7, 3)). + +%% @doc Write a block into the node's block cache at the height pseudo-path +%% `~copycat@1.0/arweave' caches it under, so the enumerator reads it locally +%% rather than from the network. +write_test_block(Height, TXs, Opts) -> + {ok, MsgID} = + hb_cache:write(#{ <<"height">> => Height, <<"txs">> => TXs }, Opts), + hb_cache:link( + MsgID, + hb_path:to_binary( + [?ARWEAVE_DEVICE, <<"block">>, <<"height">>, hb_util:bin(Height)] + ), + Opts + ), + ok. + +%% @doc In `all' mode the schedule is enumerated from the block headers, in +%% chain order, and the process's own transaction is not re-assigned: it is +%% already slot 0. +enumerate_blocks_test() -> + Store = hb_test_utils:test_store(), + hb_store:start(Store), + Opts = #{ <<"store">> => [Store] }, + ProcID = hb_util:human_id(crypto:strong_rand_bytes(32)), + First = hb_util:human_id(crypto:strong_rand_bytes(32)), + Second = hb_util:human_id(crypto:strong_rand_bytes(32)), + Third = hb_util:human_id(crypto:strong_rand_bytes(32)), + ok = write_test_block(10, [First, ProcID], Opts), + ok = write_test_block(11, [Second, Third], Opts), + ?assertEqual( + {ok, [{10, First}, {11, Second}, {11, Third}]}, + enumerate_blocks(ProcID, 10, 11, Opts) + ). + +%% @doc Block-enumerated transactions keep chain order rather than being sorted +%% by offset -- offsets tie for every transaction that carries no data, which +%% is what an AO message is -- and each records the height that sequenced it. +%% Bundled data items are still dropped. +base_layer_blocks_test() -> + Store = hb_test_utils:test_store(), + hb_store:start(Store), + ArwStore = #{ <<"index-store">> => [Store] }, + Opts = #{ <<"arweave-index-store">> => ArwStore }, + Early = hb_util:human_id(crypto:strong_rand_bytes(32)), + Late = hb_util:human_id(crypto:strong_rand_bytes(32)), + Bundled = hb_util:human_id(crypto:strong_rand_bytes(32)), + ok = hb_store_arweave:write_offset(ArwStore, Late, <<"tx@1.0">>, 100, 0), + ok = hb_store_arweave:write_offset(ArwStore, Bundled, <<"ans104@1.0">>, 150, 0), + ok = hb_store_arweave:write_offset(ArwStore, Early, <<"tx@1.0">>, 200, 0), + ?assertEqual( + {ok, + [ + {#{ <<"offset">> => 100, <<"block-height">> => 10 }, Late}, + {#{ <<"offset">> => 200, <<"block-height">> => 11 }, Early} + ] + }, + base_layer_blocks([{10, Late}, {11, Bundled}, {11, Early}], Opts) + ). + +%% @doc A transaction the block lists but the index does not hold means the +%% range was not fully indexed. Dropping it would shorten the schedule and, as +%% slots are positional, shift every later slot on this node alone -- so the +%% range fails and `synced-to' stays where it was. +base_layer_blocks_unindexed_test() -> + Store = hb_test_utils:test_store(), + hb_store:start(Store), + ArwStore = #{ <<"index-store">> => [Store] }, + Opts = #{ <<"arweave-index-store">> => ArwStore }, + Indexed = hb_util:human_id(crypto:strong_rand_bytes(32)), + Missing = hb_util:human_id(crypto:strong_rand_bytes(32)), + ok = hb_store_arweave:write_offset(ArwStore, Indexed, <<"tx@1.0">>, 100, 0), + ?assertMatch( + {error, #{ <<"status">> := 503 }}, + base_layer_blocks([{10, Indexed}, {10, Missing}], Opts) + ). + +%% @doc An `all'-mode assignment records the height that sequenced it as well +%% as its weave offset, and both survive the cache round trip that a process +%% reads its schedule back through. This is the contract +%% `~arweave-swap@1.0' reads its clock from. +all_mode_assignment_test() -> + Store = hb_test_utils:test_store(), + hb_store:start(Store), + Opts = #{ <<"store">> => [Store], <<"priv-wallet">> => ar_wallet:new() }, + ProcID = hb_util:human_id(crypto:strong_rand_bytes(32)), + Msg = + hb_message:commit( + #{ <<"target">> => ProcID }, + Opts, + #{ <<"commitment-device">> => <<"tx@1.0">> } + ), + ok = + write_assignment( + ProcID, + 1, + #{ <<"offset">> => 42, <<"block-height">> => 1958986 }, + Msg, + Opts + ), + {ok, Assignment} = dev_arweave_scheduler_cache:read(ProcID, 1, Opts), + ?assertEqual(1, hb_util:int(hb_ao:get(<<"slot">>, Assignment, Opts))), + ?assertEqual(42, hb_util:int(hb_ao:get(<<"offset">>, Assignment, Opts))), + ?assertEqual( + 1958986, + hb_util:int(hb_ao:get(<<"block-height">>, Assignment, Opts)) + ). + +%% @doc The mode is pinned in the persisted state, so a schedule can never be +%% half-derived in each mode. +state_mode_round_trip_test() -> + Store = hb_test_utils:test_store(), + hb_store:start(Store), + Opts = #{ <<"store">> => [Store] }, + ProcID = hb_util:human_id(crypto:strong_rand_bytes(32)), + State = + #{ + <<"next-slot">> => 1, + <<"spawn-height">> => 10, + <<"synced-to">> => 9, + <<"mode">> => <<"all">> + }, + ok = dev_arweave_scheduler_cache:write_state(ProcID, State, Opts), + ?assertEqual({ok, State}, dev_arweave_scheduler_cache:read_state(ProcID, Opts)). + %% @doc `/status' reports each tracked process's contiguously-indexed block %% range straight from the cache, without triggering a synchronization. status_attestable_range_test() -> @@ -893,7 +1211,8 @@ status_attestable_range_test() -> #{ <<"next-slot">> => 4, <<"spawn-height">> => 100, - <<"synced-to">> => 150 + <<"synced-to">> => 150, + <<"mode">> => ?DEFAULT_MODE }, Opts ), diff --git a/src/preloaded/process/dev_arweave_scheduler_cache.erl b/src/preloaded/process/dev_arweave_scheduler_cache.erl index a8af37424..e94e93084 100644 --- a/src/preloaded/process/dev_arweave_scheduler_cache.erl +++ b/src/preloaded/process/dev_arweave_scheduler_cache.erl @@ -37,9 +37,11 @@ list_processes(RawOpts) -> %% @doc Read the persisted synchronization state for a process. Returns %% `{ok, State}' with the `next-slot' to be assigned, the process's -%% `spawn-height', and `synced-to' (the highest block whose messages have been -%% contiguously indexed and materialized), or propagates the store's -%% `not_found'. +%% `spawn-height', `synced-to' (the highest block whose messages have been +%% contiguously indexed and materialized), and the `mode' the process is +%% sequenced in, or propagates the store's `not_found'. The mode is pinned here +%% at first contact: it is read from the process message, which cannot change, +%% and a schedule half-derived in each mode would be undetectable. read_state(ProcID, RawOpts) -> Opts = opts(RawOpts), Store = hb_opts:get(store, no_viable_store, Opts), @@ -50,11 +52,14 @@ read_state(ProcID, RawOpts) -> hb_store:read(Store, state_path(ProcID, <<"spawn-height">>), Opts), {ok, SyncedTo} ?= hb_store:read(Store, state_path(ProcID, <<"synced-to">>), Opts), + {ok, Mode} ?= + hb_store:read(Store, state_path(ProcID, <<"mode">>), Opts), {ok, #{ <<"next-slot">> => hb_util:int(NextSlot), <<"spawn-height">> => hb_util:int(SpawnHeight), - <<"synced-to">> => hb_util:int(SyncedTo) + <<"synced-to">> => hb_util:int(SyncedTo), + <<"mode">> => Mode } } end. @@ -66,14 +71,16 @@ write_state(ProcID, State, RawOpts) -> #{ <<"next-slot">> := NextSlot, <<"spawn-height">> := SpawnHeight, - <<"synced-to">> := SyncedTo + <<"synced-to">> := SyncedTo, + <<"mode">> := Mode } = State, hb_store:write( Store, #{ state_path(ProcID, <<"next-slot">>) => hb_util:bin(NextSlot), state_path(ProcID, <<"spawn-height">>) => hb_util:bin(SpawnHeight), - state_path(ProcID, <<"synced-to">>) => hb_util:bin(SyncedTo) + state_path(ProcID, <<"synced-to">>) => hb_util:bin(SyncedTo), + state_path(ProcID, <<"mode">>) => Mode }, Opts ). From 60ae372f61e0f560900bc12fda0945e0d409ecf4 Mon Sep 17 00:00:00 2001 From: Sam Williams Date: Sat, 25 Jul 2026 00:47:29 -0400 Subject: [PATCH 7/7] fix(arweave-scheduler): read a header without reaching for its data The device promises that assignment bodies are transaction headers only, so that a schedule never depends on the availability of anybody's data. Reading one through the whole store chain broke that promise: the read falls through to a gateway, which answers with the entire transaction -- its data, or, for a bundle, its items decoded into submessages whose contents this node does not hold -- and `lib_scheduler:write_assignment' then forces the assignment to load. The slot dies on a link to a chunk nobody fetched. In `target' mode this never came up, because the only messages in a schedule were the header-only ones addressed to the process. In `all' mode every transaction on the network is a message and most of them carry data, so it came up immediately: replaying a five-block range failed on a stranger's bundle. Headers are now read from the node's own stores, where the indexing pass has cached a data-free header for every plain transaction it walked, and anything else is fetched as a header in its own right with `exclude-data'. Loading whatever happened to be available locally would have been worse than either: two nodes would mint different assignments for the same transaction. Found by replaying a real sale off mainnet; the live-weave fixture tests for `target' mode are unchanged and still pass. --- .../process/dev_arweave_scheduler.erl | 19 ++++++++++++++++--- 1 file changed, 16 insertions(+), 3 deletions(-) diff --git a/src/preloaded/process/dev_arweave_scheduler.erl b/src/preloaded/process/dev_arweave_scheduler.erl index 3977ef3a0..93c5d0efb 100644 --- a/src/preloaded/process/dev_arweave_scheduler.erl +++ b/src/preloaded/process/dev_arweave_scheduler.erl @@ -832,9 +832,22 @@ slot_zero(_Mode, Offset, _SpawnHeight) -> %% tx@1.0 signature -- so its committed ID is the transaction ID and it verifies %% -- and no data is attached at this layer, so the schedule never depends on %% the availability of any transaction's data. +%% The header is taken as it is read, links and all. Forcing it to load in full +%% would demand a transaction's data -- the very thing this schedule promises +%% never to depend on -- and any transaction on the network can be a message +%% here, so most of them carry some. Loading only what happens to be available +%% locally would be worse than either: two nodes would mint different +%% assignments for the same transaction. +%% Only the node's own stores are consulted, because only they are known to +%% hold what this schedule wants: the indexing pass caches a data-free header +%% for every plain transaction it walks. Letting the read fall through to a +%% gateway would answer with the whole transaction instead -- data, or a bundle +%% decoded into items whose contents are not on this node -- and the schedule +%% would then depend on the availability of data it promises never to touch. +%% Anything not cached locally is fetched as a header in its own right. read_tx_header(TXID, Opts) -> - case hb_cache:read(TXID, Opts) of - {ok, Msg} -> {ok, hb_cache:ensure_all_loaded(Msg, Opts)}; + case hb_cache:read(TXID, hb_store:scope(Opts, local)) of + {ok, Msg} -> {ok, Msg}; _ -> fetch_tx_header(TXID, Opts) end. @@ -855,7 +868,7 @@ fetch_tx_header(TXID, Opts) -> no_result_cache(Opts) ), case Res of - {ok, Msg} -> {ok, hb_cache:ensure_all_loaded(Msg, Opts)}; + {ok, Msg} -> {ok, Msg}; _ -> {error, #{