Current section
Files
Jump to
Current section
Files
lib/mongo/monitor.ex
# see https://github.com/mongodb/specifications/blob/master/source/server-discovery-and-monitoring/server-discovery-and-monitoring.rst#monitoring
defmodule Mongo.Monitor do
@moduledoc false
use GenServer
use Bitwise
require Logger
alias Mongo.ServerDescription
alias Mongo.Events.ServerHeartbeatStartedEvent
alias Mongo.Events.{ServerHeartbeatStartedEvent, ServerHeartbeatFailedEvent,
ServerHeartbeatSucceededEvent}
# this is not configurable because the specification says so
# see https://github.com/mongodb/specifications/blob/master/source/server-discovery-and-monitoring/server-discovery-and-monitoring.rst#minheartbeatfrequencyms
@min_heartbeat_frequency_ms 500
def start_link(args, gen_server_opts \\ []) do
GenServer.start_link(__MODULE__, args, gen_server_opts)
end
# We need to stop asynchronously because a Monitor can call the Topology
# which may try to stop the same Monitor that called it. Ending in a timeout.
# See issues #139 for some information.
def stop(pid) do
GenServer.cast(pid, :stop)
end
def force_check(pid) do
GenServer.call(pid, :check, :infinity)
end
## GenServer callbacks
@doc false
def init([server_description, topology_pid, heartbeat_frequency_ms, connection_opts]) do
opts = # monitors don't authenticate and use the "admin" database
connection_opts
|> Keyword.put(:database, "admin")
|> Keyword.put(:skip_auth, true)
|> Keyword.put(:after_connect, {__MODULE__, :connected, [self(), topology_pid]})
|> Keyword.put(:backoff_min, heartbeat_frequency_ms)
|> Keyword.put(:backoff_max, heartbeat_frequency_ms)
|> Keyword.put(:backoff_type, :rand)
|> Keyword.put(:connection_type, :monitor)
|> Keyword.put(:topology_pid, topology_pid)
|> Keyword.put(:pool_size, 1)
|> Keyword.put(:idle_interval, 5_000)
{:ok, pid} = DBConnection.start_link(Mongo.MongoDBConnection, opts)
:ok = GenServer.cast(self(), :check)
{:ok, %{
connection_pid: pid,
topology_pid: topology_pid,
server_description: server_description,
heartbeat_frequency_ms: heartbeat_frequency_ms,
opts: opts
}}
end
@doc false
def terminate(reason, state) do
GenServer.stop(state.connection_pid, reason)
end
@doc false
def connected(_connection, me, topology_pid) do
GenServer.cast(topology_pid, {:connected, me})
end
@doc false
def handle_cast(:check, state) do
check(state)
end
def handle_cast(:stop, state) do
exit(:normal)
{:noreply, state}
end
@doc false
def handle_call(:check, _from, state) do
{_, state, diff} = check(state)
{:reply, diff, state}
end
@doc false
def handle_info(:timeout, state) do
check(state)
end
## Private functions
defp check(state) do
diff = :os.system_time(:milli_seconds) - state.server_description.last_update_time
if diff < @min_heartbeat_frequency_ms do
{:noreply, state, diff}
else
server_description = is_master(state.connection_pid, state.server_description, state.opts)
:ok = GenServer.cast(state.topology_pid, {:server_description, server_description})
{:noreply, %{state | server_description: server_description}, state.heartbeat_frequency_ms}
end
end
# TODO: Remove this try/rescue once a new version of db_connection is released
defp call_is_master(conn_pid, opts) do
start_time = System.monotonic_time
result = try do
Mongo.direct_command(conn_pid, [isMaster: 1], opts)
rescue
e -> {:error, e}
end
finish_time = System.monotonic_time
rtt = System.convert_time_unit(finish_time - start_time, :native, :millisecond)
finish_time = System.convert_time_unit(finish_time, :native, :millisecond)
{result, finish_time, rtt}
end
# see https://github.com/mongodb/specifications/blob/master/source/server-discovery-and-monitoring/server-discovery-and-monitoring.rst#network-error-when-calling-ismaster
defp is_master(conn_pid, last_server_description, opts) do
:ok = Mongo.Events.notify(%ServerHeartbeatStartedEvent{
connection_pid: conn_pid
})
{result, finish_time, rtt} = call_is_master(conn_pid, opts)
case result do
{:ok, is_master_reply} ->
notify_success(rtt, is_master_reply, conn_pid)
ServerDescription.from_is_master(last_server_description, rtt, finish_time, is_master_reply)
{:error, error} ->
if last_server_description.type in [:unknown, :possible_primary] do
notify_error(rtt, error, conn_pid)
ServerDescription.from_is_master_error(last_server_description, error)
else
{result, finish_time, rtt} = call_is_master(conn_pid, opts)
case result do
{:ok, is_master_reply} ->
notify_success(rtt, is_master_reply, conn_pid)
ServerDescription.from_is_master(last_server_description, rtt, finish_time, is_master_reply)
{:error, error} ->
notify_error(rtt, error, conn_pid)
ServerDescription.from_is_master_error(last_server_description, error)
end
end
end
end
defp notify_error(rtt, error, conn_pid) do
:ok = Mongo.Events.notify(%ServerHeartbeatFailedEvent{
duration: rtt,
failure: error,
connection_pid: conn_pid
})
end
defp notify_success(rtt, reply, conn_pid) do
:ok = Mongo.Events.notify(%ServerHeartbeatSucceededEvent{
duration: rtt,
reply: reply,
connection_pid: conn_pid
})
end
end