Packages

amqp_client

3.6.13
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
amqp_client src amqp_rpc_client.erl
Raw

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 = dict:new(),
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 exisiting RPC client.
stop(Pid) ->
gen_server:call(Pid, stop, infinity).
%% @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}, infinity).
%%--------------------------------------------------------------------------
%% 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 = dict:store(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 = dict:fetch(Id, Conts),
gen_server:reply(From, Payload),
amqp_channel:call(Channel, #'basic.ack'{delivery_tag = DeliveryTag}),
{noreply, State#state{continuations = dict:erase(Id, Conts) }}.
%% @private
code_change(_OldVsn, State, _Extra) ->
{ok, State}.