Packages
macula
3.5.0
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
Current section
Files
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. Phase 1 only supports the LOCAL dispatch path: the
%%% client-side and server-side macula_stream processes live in the
%%% same BEAM and are paired with pair/2. Phase 2 swaps the local
%%% peer for a QUIC stream owned by macula_mesh_client; the public
%%% API stays unchanged.
%%%
%%% See PLAN_MACULA_STREAMING.md (macula-architecture/plans).
%%% @end
%%%-------------------------------------------------------------------
-module(macula_stream).
-behaviour(gen_server).
%% Public API
-export([
start_link/1,
pair/2,
attach_remote/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 (also used by macula_mesh_client when bridging
%% to QUIC in Phase 2)
-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 (Phase 2+):
%% undefined — unpaired
%% {local, Pid} — Phase 1 in-process pairing
%% {remote, Client, Sid} — Phase 2 QUIC: peer lives on another node,
%% deliveries are encoded as STREAM_* frames
%% and sent out the mesh_client's QUIC stream.
-type peer() :: undefined
| {local, pid()}
| {remote, 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 remote (QUIC) peer to this stream. Used by
%% macula_mesh_client to wire a stream to its QUIC-carrier: all
%% deliveries out of this stream will be encoded as STREAM_* frames
%% and sent through Client's relay connection, tagged with
%% StreamId. Deliveries INTO this stream still arrive via the
%% deliver_chunk/end/error/reply casts below — mesh_client decodes
%% the incoming frames and forwards them to the local stream pid it
%% tracks for this StreamId.
-spec attach_remote(pid(), pid(), stream_id()) -> ok.
attach_remote(StreamPid, ClientPid, StreamId)
when is_pid(StreamPid), is_pid(ClientPid), is_binary(StreamId) ->
gen_server:call(StreamPid, {pair_remote, ClientPid, 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_remote, ClientPid, StreamId}, _From, State) ->
_ = erlang:monitor(process, ClientPid),
{reply, ok, State#state{peer = {remote, ClientPid, 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, Client, Sid} — hand off to macula_mesh_client which
%% encodes the action as a STREAM_* frame
%% and sends it out the QUIC stream.
%%
%% 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, Client, Sid}} = S, Action) ->
send_remote(Client, Sid, Action, S#state.seq_out).
send_remote(Client, Sid, {chunk, Encoding, Body}, Seq) ->
macula_mesh_client:send_stream_frame(Client, stream_data, #{
stream_id => Sid,
seq => Seq,
body => Body,
encoding => Encoding
});
send_remote(Client, Sid, {end_stream, Role}, _Seq) ->
macula_mesh_client:send_stream_frame(Client, stream_end, #{
stream_id => Sid,
role => Role
});
send_remote(Client, Sid, {error, Code, Message}, _Seq) ->
macula_mesh_client:send_stream_frame(Client, stream_error, #{
stream_id => Sid,
code => Code,
message => Message
});
send_remote(Client, Sid, {reply, {ok, Value}}, _Seq) ->
macula_mesh_client:send_stream_frame(Client, stream_reply, #{
stream_id => Sid,
result => Value
});
send_remote(Client, Sid, {reply, {error, Reason}}, _Seq) ->
macula_mesh_client:send_stream_frame(Client, stream_reply, #{
stream_id => Sid,
error => #{code => <<"error">>,
message => iolist_to_binary(io_lib:format("~p", [Reason]))}
}).
%% @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, 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.