Packages

amqp_client

3.7.2
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_direct_connection.erl
Raw

src/amqp_direct_connection.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.
%%
%% @private
-module(amqp_direct_connection).
-include("amqp_client_internal.hrl").
-behaviour(amqp_gen_connection).
-export([server_close/3]).
-export([init/0, terminate/2, connect/4, do/2, open_channel_args/1, i/2,
info_keys/0, handle_message/2, closing/3, channels_terminated/1]).
-export([socket_adapter_info/2]).
-record(state, {node,
user,
vhost,
params,
adapter_info,
collector,
closing_reason, %% undefined | Reason
connected_at
}).
-define(INFO_KEYS, [type]).
-define(CREATION_EVENT_KEYS, [pid, protocol, host, port, name,
peer_host, peer_port,
user, vhost, client_properties, type,
connected_at, node, user_who_performed_action]).
%%---------------------------------------------------------------------------
%% amqp_connection:close() logically closes from the client end. We may
%% want to close from the server end.
server_close(ConnectionPid, Code, Text) ->
Close = #'connection.close'{reply_text = Text,
reply_code = Code,
class_id = 0,
method_id = 0},
amqp_gen_connection:server_close(ConnectionPid, Close).
init() ->
{ok, #state{}}.
open_channel_args(#state{node = Node,
user = User,
vhost = VHost,
collector = Collector}) ->
[self(), Node, User, VHost, Collector].
do(_Method, _State) ->
ok.
handle_message({force_event_refresh, Ref}, State = #state{node = Node}) ->
rpc:call(Node, rabbit_event, notify,
[connection_created, connection_info(State), Ref]),
{ok, State};
handle_message(closing_timeout, State = #state{closing_reason = Reason}) ->
{stop, {closing_timeout, Reason}, State};
handle_message({'DOWN', _MRef, process, _ConnSup, shutdown}, State) ->
{stop, {shutdown, node_down}, State};
handle_message({'DOWN', _MRef, process, _ConnSup, Reason}, State) ->
{stop, {remote_node_down, Reason}, State};
handle_message({'EXIT', Pid, Reason}, State) ->
{stop, rabbit_misc:format("stopping because dependent process ~p died: ~p", [Pid, Reason]), State};
handle_message(Msg, State) ->
{stop, {unexpected_msg, Msg}, State}.
closing(_ChannelCloseType, Reason, State) ->
{ok, State#state{closing_reason = Reason}}.
channels_terminated(State = #state{closing_reason = Reason,
collector = Collector}) ->
rabbit_queue_collector_common:delete_all(Collector),
{stop, {shutdown, Reason}, State}.
terminate(_Reason, #state{node = Node} = State) ->
rpc:call(Node, rabbit_direct, disconnect,
[self(), [{pid, self()},
{node, Node},
{name, i(name, State)}]]),
ok.
i(type, _State) -> direct;
i(pid, _State) -> self();
%% Mandatory connection parameters
i(node, #state{node = N}) -> N;
i(user, #state{params = P}) -> P#amqp_params_direct.username;
i(user_who_performed_action, St) -> i(user, St);
i(vhost, #state{params = P}) -> P#amqp_params_direct.virtual_host;
i(client_properties, #state{params = P}) ->
P#amqp_params_direct.client_properties;
i(connected_at, #state{connected_at = T}) -> T;
%%
%% Optional adapter info
%%
%% adapter_info can be undefined e.g. when we were
%% not granted access to a vhost
i(_Key, #state{adapter_info = undefined}) -> unknown;
i(protocol, #state{adapter_info = I}) -> I#amqp_adapter_info.protocol;
i(host, #state{adapter_info = I}) -> I#amqp_adapter_info.host;
i(port, #state{adapter_info = I}) -> I#amqp_adapter_info.port;
i(peer_host, #state{adapter_info = I}) -> I#amqp_adapter_info.peer_host;
i(peer_port, #state{adapter_info = I}) -> I#amqp_adapter_info.peer_port;
i(name, #state{adapter_info = I}) -> I#amqp_adapter_info.name;
i(internal_user, #state{user = U}) -> U;
i(Item, _State) -> throw({bad_argument, Item}).
info_keys() ->
?INFO_KEYS.
infos(Items, State) ->
[{Item, i(Item, State)} || Item <- Items].
connection_info(State = #state{adapter_info = I}) ->
infos(?CREATION_EVENT_KEYS, State) ++ I#amqp_adapter_info.additional_info.
connect(Params = #amqp_params_direct{username = Username,
password = Password,
node = Node,
adapter_info = Info,
virtual_host = VHost},
SIF, _TypeSup, State) ->
State1 = State#state{node = Node,
vhost = VHost,
params = Params,
adapter_info = ensure_adapter_info(Info),
connected_at =
os:system_time(milli_seconds)},
case rpc:call(Node, rabbit_direct, connect,
[{Username, Password}, VHost, ?PROTOCOL, self(),
connection_info(State1)]) of
{ok, {User, ServerProperties}} ->
{ok, ChMgr, Collector} = SIF(i(name, State1)),
State2 = State1#state{user = User,
collector = Collector},
%% There's no real connection-level process on the remote
%% node for us to monitor or link to, but we want to
%% detect connection death if the remote node goes down
%% when there are no channels. So we monitor the
%% supervisor; that way we find out if the node goes down
%% or the rabbit app stops.
erlang:monitor(process, {rabbit_direct_client_sup, Node}),
{ok, {ServerProperties, 0, ChMgr, State2}};
{error, _} = E ->
E;
{badrpc, nodedown} ->
{error, {nodedown, Node}}
end.
ensure_adapter_info(none) ->
ensure_adapter_info(#amqp_adapter_info{});
ensure_adapter_info(A = #amqp_adapter_info{protocol = unknown}) ->
ensure_adapter_info(A#amqp_adapter_info{
protocol = {'Direct', ?PROTOCOL:version()}});
ensure_adapter_info(A = #amqp_adapter_info{name = unknown}) ->
Name = list_to_binary(rabbit_misc:pid_to_string(self())),
ensure_adapter_info(A#amqp_adapter_info{name = Name});
ensure_adapter_info(Info) -> Info.
socket_adapter_info(Sock, Protocol) ->
{PeerHost, PeerPort, Host, Port} =
case rabbit_net:socket_ends(Sock, inbound) of
{ok, Res} -> Res;
_ -> {unknown, unknown, unknown, unknown}
end,
Name = case rabbit_net:connection_string(Sock, inbound) of
{ok, Res1} -> Res1;
_Error -> "(unknown)"
end,
#amqp_adapter_info{protocol = Protocol,
name = list_to_binary(Name),
host = Host,
port = Port,
peer_host = PeerHost,
peer_port = PeerPort,
additional_info = maybe_ssl_info(Sock)}.
maybe_ssl_info(Sock) ->
RealSocket = rabbit_net:unwrap_socket(Sock),
case rabbit_net:is_ssl(RealSocket) of
true -> [{ssl, true}] ++ ssl_info(RealSocket) ++ ssl_cert_info(RealSocket);
false -> [{ssl, false}]
end.
ssl_info(Sock) ->
{Protocol, KeyExchange, Cipher, Hash} =
case rabbit_net:ssl_info(Sock) of
{ok, Infos} ->
{_, P} = lists:keyfind(protocol, 1, Infos),
case lists:keyfind(cipher_suite, 1, Infos) of
{_,{K, C, H}} -> {P, K, C, H};
{_,{K, C, H, _}} -> {P, K, C, H}
end;
_ ->
{unknown, unknown, unknown, unknown}
end,
[{ssl_protocol, Protocol},
{ssl_key_exchange, KeyExchange},
{ssl_cipher, Cipher},
{ssl_hash, Hash}].
ssl_cert_info(Sock) ->
case rabbit_net:peercert(Sock) of
{ok, Cert} ->
[{peer_cert_issuer, list_to_binary(
rabbit_cert_info:issuer(Cert))},
{peer_cert_subject, list_to_binary(
rabbit_cert_info:subject(Cert))},
{peer_cert_validity, list_to_binary(
rabbit_cert_info:validity(Cert))}];
_ ->
[]
end.