Packages
amqp_client
3.7.21
4.2.1
4.2.0
retired
4.1.6
4.1.5
retired
4.1.5-rc.2
retired
4.1.5-rc.1
retired
4.1.5-1
4.0.3
4.0.3-rc.1
3.13.0-rc.2
3.13.0-rc.1
3.12.14
3.12.13
3.12.12
3.12.11
3.12.10
3.12.9
3.12.8
3.12.7
3.12.6
3.12.5
3.12.4
3.12.3
3.12.2
3.12.1
3.12.0
3.12.0-rc.4
3.11.28
3.11.27
3.11.26
3.11.25
3.11.24
3.11.23
3.11.22
3.11.21
3.11.20
3.11.19
3.11.18
3.11.17
3.11.13
3.11.12
3.11.11
3.11.10
3.11.9
3.11.8
3.11.7
3.11.6
3.11.5
3.11.4
3.11.3
3.11.2
3.11.1
3.11.0
3.11.0-rc.2
3.11.0-rc.1
3.11.0-1
3.10.20
3.10.19
3.10.18
3.10.17
3.10.16
3.10.15
3.10.14
3.10.13
3.10.12
3.10.11
3.10.10
3.10.9
3.10.8
3.10.7
3.10.6
3.10.5
3.10.4
3.10.3
3.10.2
3.10.1
3.10.0
3.10.0-rc.6
3.10.0-rc.5
3.9.29
3.9.28
3.9.27
3.9.26
3.9.25
3.9.24
3.9.23
3.9.22
3.9.21
3.9.20
3.9.19
3.9.18
3.9.17
3.9.16
3.9.15
3.9.11
3.9.8
3.9.5
3.9.4
3.9.1
3.8.35
3.8.34
3.8.33
3.8.32
3.8.31
3.8.30
3.8.25
3.8.21
3.8.19
3.8.14
3.8.11
3.8.10
3.8.10-rc.1
3.8.9
3.8.7
3.8.6
3.8.6-rc.2
3.8.6-rc.1
3.8.5
3.8.5-rc.2
3.8.5-rc.1
3.8.4
3.8.4-rc.3
3.8.4-rc.1
3.8.3
3.8.3-rc.2
3.8.3-rc.1
3.8.2
3.8.2-rc.1
3.8.1
3.8.1-rc.1
3.8.0
3.8.0-rc.3
3.8.0-rc.2
3.8.0-rc.1
3.7.28
3.7.27
3.7.27-rc.2
3.7.27-rc.1
3.7.26
3.7.25
3.7.25-rc.1
3.7.24
3.7.24-rc.2
3.7.24-rc.1
3.7.23
3.7.23-rc.1
3.7.22
3.7.22-rc.2
3.7.22-rc.1
3.7.21
3.7.20
3.7.20-rc.2
3.7.20-rc.1
3.7.19
3.7.18
3.7.18-rc.1
3.7.17
3.7.17-rc.3
3.7.17-rc.2
3.7.15
3.7.14
3.7.14-rc.2
3.7.14-rc.1
3.7.13
3.7.13-rc.2
3.7.13-rc.1
3.7.12
3.7.12-rc.2
3.7.12-rc.1
3.7.11
3.7.11-rc.2
3.7.11-rc.1
3.7.10-rc.4
3.7.10-rc.3
3.7.10-rc.2
3.7.10-rc.1
3.7.9
3.7.9-rc.3
3.7.9-rc.2
3.7.8
3.7.8-rc.4
3.7.8-rc.3
3.7.8-rc.2
3.7.8-rc.1
3.7.7
3.7.7-rc.2
3.7.7-rc.1
3.7.6
3.7.6-rc.2
3.7.6-rc.1
3.7.5
3.7.5-rc.1
3.7.4
3.7.4-rc.4
3.7.4-rc.3
3.7.4-rc.2
3.7.4-rc.1
3.7.3
3.7.3-rc.2
3.7.3-rc.1
3.7.2
3.7.1
3.7.0-rc.2
3.6.16
3.6.16-rc.1
3.6.15
3.6.15-rc.1
3.6.14
3.6.13
3.6.12
3.6.11
3.6.10
3.6.9
3.6.8
3.6.7
3.6.7-pre.1
3.5.6
3.5.0
3.4.0
3.3.5
3.0.2
0.0.0-rc.1
RabbitMQ AMQP Client
Current section
Files
Jump to
Current section
Files
src/amqp_rpc_client.erl
%% The contents of this file are subject to the Mozilla Public License
%% Version 1.1 (the "License"); you may not use this file except in
%% compliance with the License. You may obtain a copy of the License at
%% http://www.mozilla.org/MPL/
%%
%% Software distributed under the License is distributed on an "AS IS"
%% basis, WITHOUT WARRANTY OF ANY KIND, either express or implied. See the
%% License for the specific language governing rights and limitations
%% under the License.
%%
%% The Original Code is RabbitMQ.
%%
%% The Initial Developer of the Original Code is GoPivotal, Inc.
%% Copyright (c) 2007-2017 Pivotal Software, Inc. All rights reserved.
%%
%% @doc This module allows the simple execution of an asynchronous RPC over
%% AMQP. It frees a client programmer of the necessary having to AMQP
%% plumbing. Note that the this module does not handle any data encoding,
%% so it is up to the caller to marshall and unmarshall message payloads
%% accordingly.
-module(amqp_rpc_client).
-include("amqp_client.hrl").
-behaviour(gen_server).
-export([start/2, start_link/2, stop/1]).
-export([call/2]).
-export([init/1, terminate/2, code_change/3, handle_call/3,
handle_cast/2, handle_info/2]).
-record(state, {channel,
reply_queue,
exchange,
routing_key,
continuations = #{},
correlation_id = 0}).
%%--------------------------------------------------------------------------
%% API
%%--------------------------------------------------------------------------
%% @spec (Connection, Queue) -> RpcClient
%% where
%% Connection = pid()
%% Queue = binary()
%% RpcClient = pid()
%% @doc Starts a new RPC client instance that sends requests to a
%% specified queue. This function returns the pid of the RPC client process
%% that can be used to invoke RPCs and stop the client.
start(Connection, Queue) ->
{ok, Pid} = gen_server:start(?MODULE, [Connection, Queue], []),
Pid.
%% @spec (Connection, Queue) -> RpcClient
%% where
%% Connection = pid()
%% Queue = binary()
%% RpcClient = pid()
%% @doc Starts, and links to, a new RPC client instance that sends requests
%% to a specified queue. This function returns the pid of the RPC client
%% process that can be used to invoke RPCs and stop the client.
start_link(Connection, Queue) ->
{ok, Pid} = gen_server:start_link(?MODULE, [Connection, Queue], []),
Pid.
%% @spec (RpcClient) -> ok
%% where
%% RpcClient = pid()
%% @doc Stops an existing RPC client.
stop(Pid) ->
gen_server:call(Pid, stop, amqp_util:call_timeout()).
%% @spec (RpcClient, Payload) -> ok
%% where
%% RpcClient = pid()
%% Payload = binary()
%% @doc Invokes an RPC. Note the caller of this function is responsible for
%% encoding the request and decoding the response.
call(RpcClient, Payload) ->
gen_server:call(RpcClient, {call, Payload}, amqp_util:call_timeout()).
%%--------------------------------------------------------------------------
%% Plumbing
%%--------------------------------------------------------------------------
%% Sets up a reply queue for this client to listen on
setup_reply_queue(State = #state{channel = Channel}) ->
#'queue.declare_ok'{queue = Q} =
amqp_channel:call(Channel, #'queue.declare'{exclusive = true,
auto_delete = true}),
State#state{reply_queue = Q}.
%% Registers this RPC client instance as a consumer to handle rpc responses
setup_consumer(#state{channel = Channel, reply_queue = Q}) ->
amqp_channel:call(Channel, #'basic.consume'{queue = Q}).
%% Publishes to the broker, stores the From address against
%% the correlation id and increments the correlationid for
%% the next request
publish(Payload, From,
State = #state{channel = Channel,
reply_queue = Q,
exchange = X,
routing_key = RoutingKey,
correlation_id = CorrelationId,
continuations = Continuations}) ->
EncodedCorrelationId = base64:encode(<<CorrelationId:64>>),
Props = #'P_basic'{correlation_id = EncodedCorrelationId,
content_type = <<"application/octet-stream">>,
reply_to = Q},
Publish = #'basic.publish'{exchange = X,
routing_key = RoutingKey,
mandatory = true},
amqp_channel:call(Channel, Publish, #amqp_msg{props = Props,
payload = Payload}),
State#state{correlation_id = CorrelationId + 1,
continuations = maps:put(EncodedCorrelationId, From, Continuations)}.
%%--------------------------------------------------------------------------
%% gen_server callbacks
%%--------------------------------------------------------------------------
%% Sets up a reply queue and consumer within an existing channel
%% @private
init([Connection, RoutingKey]) ->
{ok, Channel} = amqp_connection:open_channel(
Connection, {amqp_direct_consumer, [self()]}),
InitialState = #state{channel = Channel,
exchange = <<>>,
routing_key = RoutingKey},
State = setup_reply_queue(InitialState),
setup_consumer(State),
{ok, State}.
%% Closes the channel this gen_server instance started
%% @private
terminate(_Reason, #state{channel = Channel}) ->
amqp_channel:close(Channel),
ok.
%% Handle the application initiated stop by just stopping this gen server
%% @private
handle_call(stop, _From, State) ->
{stop, normal, ok, State};
%% @private
handle_call({call, Payload}, From, State) ->
NewState = publish(Payload, From, State),
{noreply, NewState}.
%% @private
handle_cast(_Msg, State) ->
{noreply, State}.
%% @private
handle_info({#'basic.consume'{}, _Pid}, State) ->
{noreply, State};
%% @private
handle_info(#'basic.consume_ok'{}, State) ->
{noreply, State};
%% @private
handle_info(#'basic.cancel'{}, State) ->
{noreply, State};
%% @private
handle_info(#'basic.cancel_ok'{}, State) ->
{stop, normal, State};
%% @private
handle_info({#'basic.deliver'{delivery_tag = DeliveryTag},
#amqp_msg{props = #'P_basic'{correlation_id = Id},
payload = Payload}},
State = #state{continuations = Conts, channel = Channel}) ->
From = maps:get(Id, Conts),
gen_server:reply(From, Payload),
amqp_channel:call(Channel, #'basic.ack'{delivery_tag = DeliveryTag}),
{noreply, State#state{continuations = maps:remove(Id, Conts) }}.
%% @private
code_change(_OldVsn, State, _Extra) ->
{ok, State}.