Packages
amqp_client
3.8.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_server.erl
%% This Source Code Form is subject to the terms of the Mozilla Public
%% License, v. 2.0. If a copy of the MPL was not distributed with this
%% file, You can obtain one at https://mozilla.org/MPL/2.0/.
%%
%% Copyright (c) 2007-2021 VMware, Inc. or its affiliates. All rights reserved.
%%
%% @doc This is a utility module that is used to expose an arbitrary function
%% via an asynchronous RPC over AMQP mechanism. It frees the implementor of
%% a simple function from having to plumb this into AMQP. Note that the
%% RPC server does not handle any data encoding, so it is up to the callback
%% function to marshall and unmarshall message payloads accordingly.
-module(amqp_rpc_server).
-behaviour(gen_server).
-include("amqp_client.hrl").
-export([init/1, terminate/2, code_change/3, handle_call/3,
handle_cast/2, handle_info/2]).
-export([start/3, start_link/3]).
-export([stop/1]).
-record(state, {channel,
handler}).
%%--------------------------------------------------------------------------
%% API
%%--------------------------------------------------------------------------
%% @spec (Connection, Queue, RpcHandler) -> RpcServer
%% where
%% Connection = pid()
%% Queue = binary()
%% RpcHandler = function()
%% RpcServer = pid()
%% @doc Starts a new RPC server instance that receives requests via a
%% specified queue and dispatches them to a specified handler function. This
%% function returns the pid of the RPC server that can be used to stop the
%% server.
start(Connection, Queue, Fun) ->
{ok, Pid} = gen_server:start(?MODULE, [Connection, Queue, Fun], []),
Pid.
%% @spec (Connection, Queue, RpcHandler) -> RpcServer
%% where
%% Connection = pid()
%% Queue = binary()
%% RpcHandler = function()
%% RpcServer = pid()
%% @doc Starts, and links to, a new RPC server instance that receives
%% requests via a specified queue and dispatches them to a specified
%% handler function. This function returns the pid of the RPC server that
%% can be used to stop the server.
start_link(Connection, Queue, Fun) ->
{ok, Pid} = gen_server:start_link(?MODULE, [Connection, Queue, Fun], []),
Pid.
%% @spec (RpcServer) -> ok
%% where
%% RpcServer = pid()
%% @doc Stops an existing RPC server.
stop(Pid) ->
gen_server:call(Pid, stop, amqp_util:call_timeout()).
%%--------------------------------------------------------------------------
%% gen_server callbacks
%%--------------------------------------------------------------------------
%% @private
init([Connection, Q, Fun]) ->
{ok, Channel} = amqp_connection:open_channel(
Connection, {amqp_direct_consumer, [self()]}),
amqp_channel:call(Channel, #'queue.declare'{queue = Q}),
amqp_channel:call(Channel, #'basic.consume'{queue = Q}),
{ok, #state{channel = Channel, handler = Fun} }.
%% @private
handle_info(shutdown, State) ->
{stop, normal, State};
%% @private
handle_info({#'basic.consume'{}, _}, 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 = Props, payload = Payload}},
State = #state{handler = Fun, channel = Channel}) ->
#'P_basic'{correlation_id = CorrelationId,
reply_to = Q} = Props,
Response = Fun(Payload),
Properties = #'P_basic'{correlation_id = CorrelationId},
Publish = #'basic.publish'{exchange = <<>>,
routing_key = Q,
mandatory = true},
amqp_channel:call(Channel, Publish, #amqp_msg{props = Properties,
payload = Response}),
amqp_channel:call(Channel, #'basic.ack'{delivery_tag = DeliveryTag}),
{noreply, State};
%% @private
handle_info({'DOWN', _MRef, process, _Pid, _Info}, State) ->
{noreply, State}.
%% @private
handle_call(stop, _From, State) ->
{stop, normal, ok, State}.
%%--------------------------------------------------------------------------
%% Rest of the gen_server callbacks
%%--------------------------------------------------------------------------
%% @private
handle_cast(_Message, State) ->
{noreply, State}.
%% Closes the channel this gen_server instance started
%% @private
terminate(_Reason, #state{channel = Channel}) ->
amqp_channel:close(Channel),
ok.
%% @private
code_change(_OldVsn, State, _Extra) ->
{ok, State}.