Packages

macula

4.4.9
7.0.0 6.0.0 5.2.2 5.2.1 5.2.0 5.1.0 5.0.0 4.8.0 4.7.1 4.7.0 4.6.0 4.5.0 4.4.10 4.4.9 4.4.8 4.4.7 4.4.6 4.4.5 4.4.4 4.4.3 4.4.2 4.4.1 4.4.0 4.3.1 4.3.0 4.2.9 4.2.8 4.2.7 4.2.6 4.2.5 4.2.4 4.2.3 4.2.2 4.2.1 4.2.0 4.1.1 4.1.0 4.0.0 3.16.0 3.15.3 3.15.2 3.15.1 3.14.0 3.13.0 3.12.1 3.12.0 3.11.1 3.11.0 3.10.3 3.10.2 3.10.1 3.9.0 3.8.0 3.7.0 3.5.0 3.4.0 3.3.0 3.2.0 3.1.0 3.0.0 2.1.1 2.1.0 2.0.0 1.5.2 1.5.1 1.4.30 1.4.29 1.4.28 1.4.27 1.4.26 1.4.25 1.4.24 1.4.23 1.4.22 1.4.21 1.4.20 1.4.19 1.4.18 1.4.17 1.4.16 1.4.15 1.4.14 1.4.13 1.4.11 1.4.10 1.4.9 1.4.8 1.4.7 1.4.6 1.4.5 1.4.4 1.4.3 1.4.2 1.4.1 1.4.0 1.3.1 1.3.0 1.2.0 1.1.0 1.0.10 1.0.9 1.0.8 1.0.7 1.0.6 1.0.5 1.0.4 1.0.3 1.0.2 1.0.1 1.0.0 0.48.6 0.48.5 0.48.4 0.48.3 0.48.2 0.48.1 0.48.0 0.47.1 0.47.0 0.46.3 0.46.1 0.46.0 0.45.3 0.45.2 0.45.1 0.45.0 0.44.2 0.44.1 0.44.0 0.43.3 0.43.2 0.43.1 0.43.0 0.42.9 0.42.8 0.42.7 0.42.6 0.42.5 0.42.4 0.42.3 0.42.2 0.42.1 0.42.0 0.41.1 0.41.0 0.40.1 0.40.0 0.39.9 0.39.8 0.39.7 0.39.6 0.39.5 0.39.4 0.39.3 0.39.2 0.39.1 0.39.0 0.38.8 0.38.7 0.38.6 0.38.5 0.38.4 0.38.3 0.38.2 0.38.1 0.38.0 0.37.7 0.37.6 0.37.5 0.37.4 0.37.3 0.37.2 0.37.1 0.37.0 0.36.6 0.36.5 0.36.4 0.36.3 0.36.2 0.36.1 0.36.0 0.35.4 0.35.3 0.35.2 0.35.1 0.35.0 0.34.1 0.34.0 0.33.1 0.33.0 0.32.5 0.32.4 0.32.3 0.32.2 0.32.1 0.32.0 0.31.9 0.31.8 0.31.7 0.31.6 0.31.5 0.31.4 0.31.3 0.31.2 0.31.1 0.31.0 0.30.10 0.30.9 0.30.8 0.30.7 0.30.6 0.30.5 0.30.4 0.30.3 0.30.2 0.30.1 0.30.0 0.29.0 0.28.3 0.28.2 0.28.1 0.28.0 0.27.1 0.27.0 0.26.1 0.26.0 0.25.6 0.25.5 0.25.4 0.25.3 0.25.2 0.25.1 0.25.0 0.24.6 0.24.5 0.24.4 0.24.3 0.24.2 0.24.1 0.24.0 0.23.3 0.23.2 0.23.1 0.23.0 0.22.12 0.22.11 0.22.10 0.22.9 0.22.8 0.22.7 0.22.6 0.22.5 0.22.4 0.22.3 0.22.2 0.22.1 0.22.0 0.21.7 0.21.6 0.21.5 0.21.4 0.21.2 0.21.1 0.21.0 0.20.25 0.20.24 0.20.23 0.20.22 0.20.21 0.20.20 0.20.19 0.20.18 0.20.17 0.20.16 0.20.15 0.20.14 0.20.13 0.20.12 0.20.11 0.20.10 0.20.9 0.20.8 0.20.7 0.20.6 0.20.5 0.20.3 0.20.2 0.20.1 0.20.0 0.19.2 0.19.1 0.19.0 0.18.1 0.18.0 0.17.4 0.17.3 0.17.2 0.17.1 0.17.0 0.16.6 0.16.5 0.16.4 0.16.3 0.16.2 0.16.1 0.16.0 0.15.1 0.15.0 0.14.3 0.14.2 0.14.1 0.14.0 0.12.6 0.12.5 0.12.3 0.11.3 0.10.2 0.10.1 0.10.0 0.9.2 0.9.1 0.9.0 0.8.25 0.8.24 0.8.23 0.8.22 0.8.21 0.8.20 0.8.19 0.8.18 0.8.17 0.8.16 0.8.15 0.8.14 0.8.13 0.8.12 0.8.11 0.8.10 0.8.9 0.8.8 0.8.7 0.8.6 0.8.5 0.8.4 0.8.3 0.8.2 0.8.1 0.8.0 0.7.30 0.7.29 0.7.28 0.7.27 0.7.26 0.7.25 0.7.24 0.7.23 0.7.22 0.7.21 0.7.20 0.7.19 0.7.18 0.7.17 0.7.16 0.7.15 0.7.14 0.7.13 0.7.12 0.7.11 0.7.10 0.7.9 0.7.8 0.7.7 0.7.6 0.7.5 0.7.4 0.7.3 0.7.2 0.7.1 0.7.0 0.6.7 0.6.6 0.6.5 0.6.4 0.6.3 0.6.2 0.6.1 0.6.0 0.5.0 0.4.4 0.4.3 0.4.2 0.4.1 0.4.0 0.3.4 0.3.3 0.3.2 0.3.1

Macula HTTP/3 Mesh SDK — connect, subscribe, publish, call, advertise

Current section

Files

Jump to
macula src macula_stream.erl
Raw

src/macula_stream.erl

%%%-------------------------------------------------------------------
%%% @doc Macula streaming RPC — single-stream state machine.
%%%
%%% Owns one streaming RPC's state. Each `call_stream', `open_stream',
%%% or server-side handler invocation gets its own `macula_stream'
%%% gen_server. The state machine itself is carrier-agnostic; the
%%% peer shape (`{local, _}' or `{remote_via_link, _, _}') decides
%%% how chunks reach the wire.
%%%
%%% Two carriers route through `forward_to_peer/2':
%%% <ul>
%%% <li>`{local, Pid}' — in-process pairing for unit tests and
%%% `macula_stream_local' dispatch.</li>
%%% <li>`{remote_via_link, Link, Sid}' — V2 wire format via
%%% `macula_station_link' (CBOR `macula_frame:stream_*' frames
%%% over a peering connection).</li>
%%% </ul>
%%%
%%% Renamed from `macula_stream_v1' in 3.17.0; the V1 mesh_client
%%% carrier (`{remote, _, _}') was retired alongside the rest of the
%%% V1 surface in the same release. The module now spans the LOCAL
%%% carrier and the V2 station_link carrier only.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_stream).
-behaviour(gen_server).
%% Public API
-export([
start_link/1,
pair/2,
attach_to_link/3,
send/2,
send/3,
recv/1,
recv/2,
close/1,
close_send/1,
await_reply/1,
await_reply/2,
set_reply/2,
set_error/2,
abort/3,
info/1
]).
%% Peer-to-peer protocol — drives inbound deliveries from carrier
%% modules (`macula_stream_local' for LOCAL pairs, `macula_station_link'
%% for V2 station-link pairs).
-export([
deliver_chunk/3,
deliver_end/2,
deliver_error/3,
deliver_reply/2
]).
%% gen_server callbacks
-export([
init/1,
handle_call/3,
handle_cast/2,
handle_info/2,
terminate/2
]).
-type role() :: client | server.
-type mode() :: server_stream | client_stream | bidi.
-type encoding() :: raw | msgpack.
-type chunk() :: binary() | {raw, binary()} | {term, term()}.
-type stream_id() :: binary().
-type result() :: {ok, term()} | {error, term()}.
%% Peer shape:
%% undefined — unpaired
%% {local, Pid} — in-process pairing
%% {remote_via_link, L, Sid}— V2 station_link carrier: deliveries
%% encoded as `macula_frame:stream_*'
%% frames and shipped through the
%% station_link's peering connection.
-type peer() :: undefined
| {local, pid()}
| {remote_via_link, pid(), stream_id()}.
-export_type([role/0, mode/0, encoding/0, chunk/0, stream_id/0, result/0,
peer/0]).
-record(state, {
id :: stream_id(),
role :: role(),
mode :: mode(),
owner :: pid(),
owner_ref :: reference(),
peer :: peer(),
%% Recv side: inbound chunks queued, waiting recv/2 callers, eof flag
inbox = queue:new() :: queue:queue({encoding(), term()}),
waiters = queue:new() :: queue:queue({{pid(), reference()}, reference()}),
closed_recv = false :: boolean(),
%% Send side
closed_send = false :: boolean(),
seq_out = 0 :: non_neg_integer(),
seq_in = 0 :: non_neg_integer(),
%% Terminal reply (for client-stream / bidi)
reply = undefined :: undefined | result(),
reply_waiters = [] :: [{pid(), reference()}]
}).
%%%===================================================================
%%% Public API
%%%===================================================================
%% @doc Start a stream gen_server.
%%
%% Required opts: id, role, mode, owner.
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Opts) ->
gen_server:start_link(?MODULE, Opts, []).
%% @doc Pair two stream processes as peers (Phase 1 local dispatch).
-spec pair(pid(), pid()) -> ok.
pair(A, B) when is_pid(A), is_pid(B) ->
ok = gen_server:call(A, {pair_local, B}),
ok = gen_server:call(B, {pair_local, A}),
ok.
%% @doc Attach a V2 `macula_station_link' peer to this stream. The
%% station_link carries deliveries as V2 `macula_frame:stream_*'
%% frames over its peering connection (one per pool seed); inbound
%% STREAM_* frames are decoded by the link and forwarded into this
%% stream via the deliver_chunk / end / error / reply casts below.
-spec attach_to_link(pid(), pid(), stream_id()) -> ok.
attach_to_link(StreamPid, LinkPid, StreamId)
when is_pid(StreamPid), is_pid(LinkPid), is_binary(StreamId) ->
gen_server:call(StreamPid, {pair_via_link, LinkPid, StreamId}).
%% @doc Send a binary chunk on the stream.
-spec send(pid(), binary()) -> ok | {error, term()}.
send(Pid, Bin) when is_binary(Bin) ->
send(Pid, Bin, raw).
-spec send(pid(), binary() | term(), encoding()) -> ok | {error, term()}.
send(Pid, Body, raw) when is_binary(Body) ->
gen_server:call(Pid, {send, raw, Body});
send(Pid, Body, msgpack) ->
gen_server:call(Pid, {send, msgpack, Body}).
%% @doc Receive the next chunk (blocks indefinitely).
-spec recv(pid()) -> {chunk, binary()}
| {data, term()}
| eof
| {error, term()}.
recv(Pid) ->
recv(Pid, infinity).
-spec recv(pid(), timeout()) -> {chunk, binary()}
| {data, term()}
| eof
| {error, term()}.
recv(Pid, Timeout) ->
%% Long timeouts allowed because the wait is on inbound network
%% data, not on the gen_server's processing time.
GsTimeout = case Timeout of
infinity -> infinity;
N when is_integer(N) -> N + 100
end,
gen_server:call(Pid, {recv, Timeout}, GsTimeout).
%% @doc Half-close the write side. Recv side stays open.
-spec close_send(pid()) -> ok.
close_send(Pid) ->
gen_server:call(Pid, close_send).
%% @doc Close both sides. Idempotent.
-spec close(pid()) -> ok.
close(Pid) ->
gen_server:call(Pid, close).
%% @doc Wait for the terminal reply (client-stream / bidi).
-spec await_reply(pid()) -> result().
await_reply(Pid) ->
await_reply(Pid, infinity).
-spec await_reply(pid(), timeout()) -> result() | {error, timeout}.
await_reply(Pid, Timeout) ->
GsTimeout = case Timeout of
infinity -> infinity;
N when is_integer(N) -> N + 100
end,
gen_server:call(Pid, {await_reply, Timeout}, GsTimeout).
%% @doc Server-side: emit the terminal reply.
-spec set_reply(pid(), term()) -> ok.
set_reply(Pid, Result) ->
gen_server:call(Pid, {set_reply, {ok, Result}}).
%% @doc Server-side: emit a terminal error as the reply value.
-spec set_error(pid(), term()) -> ok.
set_error(Pid, Reason) ->
gen_server:call(Pid, {set_reply, {error, Reason}}).
%% @doc Abort the stream with a STREAM_ERROR frame. Both sides close;
%% any pending recv/await_reply waiters receive {error, {Code, Message}}.
-spec abort(pid(), binary(), binary()) -> ok.
abort(Pid, Code, Message) when is_binary(Code), is_binary(Message) ->
gen_server:call(Pid, {abort, Code, Message}).
%% @doc Inspect stream state (debugging).
-spec info(pid()) -> map().
info(Pid) ->
gen_server:call(Pid, info).
%%%===================================================================
%%% Peer-to-peer protocol
%%%===================================================================
%% @doc Deliver a chunk frame from the peer.
-spec deliver_chunk(pid(), encoding(), term()) -> ok.
deliver_chunk(Pid, Encoding, Body) ->
gen_server:cast(Pid, {peer_chunk, Encoding, Body}).
%% @doc Deliver a STREAM_END frame from the peer.
-spec deliver_end(pid(), send | both) -> ok.
deliver_end(Pid, Role) ->
gen_server:cast(Pid, {peer_end, Role}).
%% @doc Deliver a STREAM_ERROR frame from the peer.
-spec deliver_error(pid(), binary(), binary()) -> ok.
deliver_error(Pid, Code, Message) ->
gen_server:cast(Pid, {peer_error, Code, Message}).
%% @doc Deliver a STREAM_REPLY frame from the peer.
-spec deliver_reply(pid(), result()) -> ok.
deliver_reply(Pid, Result) ->
gen_server:cast(Pid, {peer_reply, Result}).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init(Opts) ->
Id = maps:get(id, Opts),
Role = maps:get(role, Opts),
Mode = maps:get(mode, Opts),
Owner = maps:get(owner, Opts),
OwnerRef = erlang:monitor(process, Owner),
{ok, #state{
id = Id,
role = Role,
mode = Mode,
owner = Owner,
owner_ref = OwnerRef
}}.
%% --- pair --------------------------------------------------------------
handle_call({pair_local, Peer}, _From, State) ->
_ = erlang:monitor(process, Peer),
{reply, ok, State#state{peer = {local, Peer}}};
handle_call({pair_via_link, LinkPid, StreamId}, _From, State) ->
_ = erlang:monitor(process, LinkPid),
{reply, ok, State#state{peer = {remote_via_link, LinkPid, StreamId}}};
%% --- send --------------------------------------------------------------
handle_call({send, _Encoding, _Body}, _From, #state{closed_send = true} = State) ->
{reply, {error, send_closed}, State};
handle_call({send, Encoding, Body}, _From, State) ->
case forward_to_peer(State, {chunk, Encoding, Body}) of
ok ->
{reply, ok, State#state{seq_out = State#state.seq_out + 1}};
{error, _} = Err ->
{reply, Err, State}
end;
%% --- recv --------------------------------------------------------------
handle_call({recv, Timeout}, From, State) ->
handle_recv(From, Timeout, State);
%% --- close_send --------------------------------------------------------
handle_call(close_send, _From, State) ->
State1 = case State#state.closed_send of
true -> State;
false ->
_ = forward_to_peer(State, {end_stream, send}),
State#state{closed_send = true}
end,
{reply, ok, State1};
%% --- close -------------------------------------------------------------
handle_call(close, _From, State) ->
_ = forward_to_peer(State, {end_stream, both}),
State1 = State#state{closed_send = true, closed_recv = true},
State2 = drain_waiters(eof, State1),
{reply, ok, State2};
%% --- await_reply -------------------------------------------------------
handle_call({await_reply, _Timeout}, _From, #state{reply = {ok, _} = R} = State) ->
{reply, R, State};
handle_call({await_reply, _Timeout}, _From, #state{reply = {error, _} = R} = State) ->
{reply, R, State};
handle_call({await_reply, Timeout}, From, State) ->
Ref = case Timeout of
infinity -> undefined;
N -> erlang:send_after(N, self(), {reply_timeout, From})
end,
Waiters = [{From, Ref} | State#state.reply_waiters],
{noreply, State#state{reply_waiters = Waiters}};
%% --- set_reply ---------------------------------------------------------
handle_call({set_reply, Result}, _From, State) ->
State1 = case State#state.reply of
undefined ->
_ = forward_to_peer(State, {reply, Result}),
State#state{reply = Result};
_ ->
State
end,
{reply, ok, State1};
handle_call({abort, Code, Message}, _From, State) ->
Err = {error, {Code, Message}},
_ = forward_to_peer(State, {error, Code, Message}),
State1 = State#state{closed_recv = true, closed_send = true,
reply = case State#state.reply of
undefined -> Err;
R -> R
end},
State2 = drain_waiters(Err, State1),
State3 = settle_reply_waiters_with(Err, State2),
{reply, ok, State3};
%% --- info --------------------------------------------------------------
handle_call(info, _From, State) ->
Map = #{
id => State#state.id,
role => State#state.role,
mode => State#state.mode,
peer => State#state.peer,
inbox_size => queue:len(State#state.inbox),
waiters => queue:len(State#state.waiters),
closed_recv => State#state.closed_recv,
closed_send => State#state.closed_send,
seq_out => State#state.seq_out,
seq_in => State#state.seq_in,
reply => State#state.reply
},
{reply, Map, State};
handle_call(_Msg, _From, State) ->
{reply, {error, unknown}, State}.
%% --- peer-delivered events --------------------------------------------
handle_cast({peer_chunk, _Encoding, _Body}, #state{closed_recv = true} = State) ->
{noreply, State};
handle_cast({peer_chunk, Encoding, Body}, State) ->
State1 = enqueue_or_deliver(Encoding, Body, State),
{noreply, State1#state{seq_in = State1#state.seq_in + 1}};
handle_cast({peer_end, send}, State) ->
%% Peer half-closed: no more inbound data
State1 = State#state{closed_recv = true},
State2 = drain_waiters(eof, State1),
{noreply, State2};
handle_cast({peer_end, both}, State) ->
State1 = State#state{closed_recv = true, closed_send = true},
State2 = drain_waiters(eof, State1),
State3 = settle_reply_waiters_with({error, peer_closed}, State2),
{noreply, State3};
handle_cast({peer_error, Code, Message}, State) ->
Err = {error, {Code, Message}},
State1 = State#state{closed_recv = true, closed_send = true},
State2 = drain_waiters(Err, State1),
State3 = settle_reply_waiters_with(Err, State2),
{noreply, State3};
handle_cast({peer_reply, Result}, State) ->
State1 = State#state{reply = Result},
State2 = settle_reply_waiters_with(Result, State1),
{noreply, State2};
handle_cast(_Msg, State) ->
{noreply, State}.
%% --- info / monitors / timers -----------------------------------------
handle_info({recv_timeout, From}, State) ->
%% Drop this waiter and reply timeout — only if it's still queued
{Replied, NewQ} = drop_waiter_and_reply(From, {error, timeout}, State#state.waiters),
case Replied of
true -> ok;
false -> ok % already served
end,
{noreply, State#state{waiters = NewQ}};
handle_info({reply_timeout, From}, State) ->
NewWaiters = lists:filter(
fun({F, _Ref}) when F =:= From ->
gen_server:reply(F, {error, timeout}),
false;
(_) -> true
end, State#state.reply_waiters),
{noreply, State#state{reply_waiters = NewWaiters}};
handle_info({'DOWN', Ref, process, Pid, _Reason}, State) ->
IsOwner = Ref =:= State#state.owner_ref andalso
Pid =:= State#state.owner,
handle_down(IsOwner, Pid, State);
handle_info(_Msg, State) ->
{noreply, State}.
terminate(_Reason, _State) -> ok.
%%%===================================================================
%%% Internal helpers
%%%===================================================================
%% @private Dispatch a stream-level action to the peer.
%%
%% Peer-shape-aware:
%% {local, Pid} — in-process pair; cast the symmetric
%% deliver_* helper directly.
%% {remote_via_link, L, Sid}— hand off to `macula_station_link'
%% which signs and sends a
%% `macula_frame:stream_*' frame.
%%
%% Action shapes:
%% {chunk, Encoding, Body}
%% {end_stream, send | both}
%% {error, Code, Message}
%% {reply, Result}
forward_to_peer(#state{peer = undefined}, _Action) ->
{error, no_peer};
forward_to_peer(#state{peer = {local, Pid}}, {chunk, Encoding, Body}) ->
deliver_chunk(Pid, Encoding, Body);
forward_to_peer(#state{peer = {local, Pid}}, {end_stream, Role}) ->
deliver_end(Pid, Role);
forward_to_peer(#state{peer = {local, Pid}}, {error, Code, Message}) ->
deliver_error(Pid, Code, Message);
forward_to_peer(#state{peer = {local, Pid}}, {reply, Result}) ->
deliver_reply(Pid, Result);
forward_to_peer(#state{peer = {remote_via_link, Link, Sid}} = S, Action) ->
send_via_link(Link, Sid, Action, S#state.seq_out).
%% @private V2 carrier: hand off to `macula_station_link' which signs
%% and ships a `macula_frame:stream_*' frame through its peering
%% connection. Action shapes mirror `send_remote/4'; the link
%% translates them to V2 frame specs internally.
send_via_link(Link, Sid, {chunk, Encoding, Body}, Seq) ->
macula_station_link:send_stream_frame(Link, stream_data, #{
stream_id => Sid,
seq => Seq,
encoding => Encoding,
body => Body
});
send_via_link(Link, Sid, {end_stream, Role}, _Seq) ->
macula_station_link:send_stream_frame(Link, stream_end, #{
stream_id => Sid,
role => Role
});
send_via_link(Link, Sid, {error, Code, Message}, _Seq) ->
macula_station_link:send_stream_frame(Link, stream_error, #{
stream_id => Sid,
code => Code,
message => Message
});
send_via_link(Link, Sid, {reply, {ok, Value}}, _Seq) ->
macula_station_link:send_stream_frame(Link, stream_reply, #{
stream_id => Sid,
payload => Value
});
send_via_link(Link, Sid, {reply, {error, _Reason} = Err}, _Seq) ->
macula_station_link:send_stream_frame(Link, stream_reply, #{
stream_id => Sid,
payload => Err
}).
%% @private Owner DOWN → stop. Otherwise check whether the dead pid
%% was our peer (or our peer's mesh_client for remote peers) and, if
%% so, surface as a stream error to any local readers / reply waiters.
handle_down(true, _Pid, State) ->
{stop, normal, State};
handle_down(false, Pid, #state{peer = {local, Pid}} = State) ->
propagate_peer_down(State);
handle_down(false, Pid, #state{peer = {remote_via_link, Pid, _Sid}} = State) ->
propagate_peer_down(State);
handle_down(false, _Pid, State) ->
{noreply, State}.
propagate_peer_down(State) ->
Err = {error, peer_down},
State1 = State#state{closed_recv = true, closed_send = true,
peer = undefined},
State2 = drain_waiters(Err, State1),
State3 = settle_reply_waiters_with(Err, State2),
{noreply, State3}.
%% @doc Either deliver a chunk to a waiting recv/2 caller or queue it.
enqueue_or_deliver(Encoding, Body, #state{waiters = W0} = State) ->
case queue:out(W0) of
{{value, {From, Ref}}, W1} ->
cancel_timer(Ref),
gen_server:reply(From, chunk_to_recv_result(Encoding, Body)),
State#state{waiters = W1};
{empty, _} ->
Inbox = queue:in({Encoding, Body}, State#state.inbox),
State#state{inbox = Inbox}
end.
handle_recv(From, _Timeout, #state{inbox = Inbox} = State) ->
case queue:out(Inbox) of
{{value, {Encoding, Body}}, Rest} ->
{reply, chunk_to_recv_result(Encoding, Body), State#state{inbox = Rest}};
{empty, _} when State#state.closed_recv ->
{reply, eof, State};
{empty, _} ->
queue_waiter(From, _Timeout, State)
end.
queue_waiter(From, Timeout, State) ->
Ref = case Timeout of
infinity -> undefined;
0 -> immediate;
N when is_integer(N) -> erlang:send_after(N, self(), {recv_timeout, From})
end,
case Ref of
immediate ->
{reply, {error, would_block}, State};
_ ->
Waiters = queue:in({From, Ref}, State#state.waiters),
{noreply, State#state{waiters = Waiters}}
end.
chunk_to_recv_result(raw, Body) -> {chunk, Body};
chunk_to_recv_result(msgpack, Body) -> {data, Body};
chunk_to_recv_result(Other, Body) -> {data, {Other, Body}}.
drain_waiters(Reply, State) ->
drain_waiters(Reply, State#state.waiters, State).
drain_waiters(Reply, Q, State) ->
case queue:out(Q) of
{{value, {From, Ref}}, Rest} ->
cancel_timer(Ref),
gen_server:reply(From, Reply),
drain_waiters(Reply, Rest, State#state{waiters = Rest});
{empty, _} ->
State#state{waiters = queue:new()}
end.
settle_reply_waiters_with(Result, State) ->
lists:foreach(
fun({From, Ref}) ->
cancel_timer(Ref),
gen_server:reply(From, Result)
end, State#state.reply_waiters),
State#state{reply_waiters = []}.
drop_waiter_and_reply(From, Reply, Q) ->
%% Walk the queue once, dropping the matching waiter and replying.
L = queue:to_list(Q),
{Match, Rest} = lists:partition(fun({F, _R}) -> F =:= From end, L),
case Match of
[{F, _Ref}] ->
gen_server:reply(F, Reply),
{true, queue:from_list(Rest)};
_ ->
{false, Q}
end.
cancel_timer(undefined) -> ok;
cancel_timer(immediate) -> ok;
cancel_timer(Ref) -> erlang:cancel_timer(Ref), ok.