macula_pusher behaviour (macula v13.2.1)
View SourceBehaviour for supervised content pushes (the sender side of a push-initiated upload — PLAN_PUSH_UPLOAD.md Phase 6).
start_link/5,6 returns immediately with a pid, delivers the outcome to Module:handle_pushed/2, and publishes sharing.push_started_v1 / sharing.push_completed_v1 mesh facts around the transfer — including outcome => cancelled if the pusher is stopped before the push resolves.
Unlike macula_feeder (which shares content from this node for a downloader to find and fetch later), this actively pushes bytes AT a specific, already-known recipient advertising an upload procedure (macula_upload:advertise/5,6) — the client_stream mode STREAMING_GUIDE.md names for exactly this ("an upload, a batch submit"), wrapped with macula_feeder/macula_download's own integrity machinery: macula_manifest:create/2 chunks and hashes Bytes up front, the manifest rides the stream's open-time Args (an out-of-band channel, not an in-band header chunk — see macula_manifest's from_wire/1), and chunks are sent in order over the ONE client_stream the recipient reads from.
One stream, chunks in order
Chunks go out sequentially over the one stream macula:call_stream/5 / macula_direct_dial:call_stream/5 opens: the same shape a hand written client_stream caller would use, just chunked and hashed for you.
The terminal reply
The recipient's own macula_upload verifies the reassembled bytes against the manifest (receiver-side, never sender-trusted — the sender's claimed manifest proves nothing on its own) and reports the outcome back over client_stream's own terminal-reply channel (macula_stream:set_reply/2 / set_error/2, surfaced here via macula_streamer's new handle_eof/1 callback — see that module's doc). This pusher blocks on macula:await_reply/1 for it, so handle_pushed/2 only ever sees {ok, Mcid} once the recipient has actually verified the bytes, never merely "the local send/2,3 calls all returned ok'" (which would only prove bytes were accepted onto the wire, not that they arrived correctly).
Real cancel
Holds the raw stream pid directly (a stream state field, alongside the lightweight open+send+await proxy worker that reports it back) so cancel/1 reaches it for a real, peer-visible macula_stream:abort/3 STREAM_ERROR — not a blunt local kill that leaves the recipient inferring cancellation from the connection simply going away.
Direct-dial
A push targets a specific ADVERTISED PROCEDURE (macula_upload:advertise_direct/6,7's procedure_advertisement), so start_link_direct/5,6 mirrors macula_stream_sink:start_link_direct/5,6's shape: a Procedure-based resolve, no Station parameter, reusing macula_direct_dial:call_stream/5, which resolves and dials as one step.
Stream I/O
A pusher opens its stream with call_stream/5, macula:call_stream/5 by default or macula_direct_dial:call_stream/5 for start_link_direct; sends its chunks with send/3 and close_send/1; waits for the recipient's reply with await_reply/1; aborts the stream on cancel with abort/3; and announces its facts with fact_publish, macula:publish/4 by default. start_link/7 and start_link_direct/7 take a stream_io start option, checked by macula_stream:stream_io/2, with those five stream functions and any other macula_stream:stream_io() ones, and a fact_publish start option, to run a pusher on something else, such as a test's scripted stream.
Example
-module(doc_pusher).
-behaviour(macula_pusher).
-export([init/1, handle_pushed/2]).
init(Parent) -> {ok, Parent}.
handle_pushed(Result, Parent) ->
Parent ! {pushed, Result},
{stop, normal, Parent}. {ok, Pid} = macula_pusher:start_link(doc_pusher, Pool, Realm,
<<"bulk.ingest">>, Bytes, self()).
Summary
Functions
Cancel an in-flight push. Publishes sharing.push_completed_v1 with outcome => cancelled if the push had not resolved yet.
Start a pusher. Pushes Bytes to Procedure on (Realm) via Pool.
As start_link/5, with Args passed to Module:init/1.
As start_link/6, with start options: stream_io gives the functions the pusher runs its stream on, and fact_publish the one it announces its facts with (see "Stream I/O" above).
As start_link/5, but resolves Procedure's procedure_advertisement from the DHT and dials its provider directly instead of pushing through the pool's existing links. See the "Direct-dial" section above. Requires the recipient to have advertised via macula_upload:advertise_direct/6,7, not plain advertise/5,6.
As start_link_direct/5, with Args passed to Module:init/1.
As start_link_direct/6, with start options: stream_io gives the functions the pusher runs its stream on, and fact_publish the one it announces its facts with (see "Stream I/O" above).
Types
-type start_opts() :: #{stream_io => macula_stream:stream_io(), fact_publish => macula_lifetime_announcer:publish()}.
Callbacks
Functions
-spec cancel(pid()) -> ok.
Cancel an in-flight push. Publishes sharing.push_completed_v1 with outcome => cancelled if the push had not resolved yet.
-spec start_link(module(), macula:pool(), macula:realm(), macula:procedure(), binary()) -> {ok, pid()} | {error, term()}.
Start a pusher. Pushes Bytes to Procedure on (Realm) via Pool.
-spec start_link(module(), macula:pool(), macula:realm(), macula:procedure(), binary(), term()) -> {ok, pid()} | {error, term()}.
As start_link/5, with Args passed to Module:init/1.
-spec start_link(module(), macula:pool(), macula:realm(), macula:procedure(), binary(), term(), start_opts()) -> {ok, pid()} | {error, term()}.
As start_link/6, with start options: stream_io gives the functions the pusher runs its stream on, and fact_publish the one it announces its facts with (see "Stream I/O" above).
-spec start_link_direct(module(), macula:pool(), macula:realm(), macula:procedure(), binary()) -> {ok, pid()} | {error, term()}.
As start_link/5, but resolves Procedure's procedure_advertisement from the DHT and dials its provider directly instead of pushing through the pool's existing links. See the "Direct-dial" section above. Requires the recipient to have advertised via macula_upload:advertise_direct/6,7, not plain advertise/5,6.
-spec start_link_direct(module(), macula:pool(), macula:realm(), macula:procedure(), binary(), term()) -> {ok, pid()} | {error, term()}.
As start_link_direct/5, with Args passed to Module:init/1.
-spec start_link_direct(module(), macula:pool(), macula:realm(), macula:procedure(), binary(), term(), start_opts()) -> {ok, pid()} | {error, term()}.
As start_link_direct/6, with start options: stream_io gives the functions the pusher runs its stream on, and fact_publish the one it announces its facts with (see "Stream I/O" above).