defmodule GenMetrics.GenServer.Monitor do use GenServer alias GenMetrics.GenServer.Manager alias GenMetrics.GenServer.Monitor alias GenMetrics.GenServer.Cluster alias GenMetrics.GenServer.Window alias GenMetrics.Reporter alias GenMetrics.Utils.Runtime @moduledoc false @call_cast_info [:handle_call, :handle_cast, :handle_info] @window_interval_default 1000 defstruct cluster: %Cluster{}, metrics: nil, start: 0, duration: 0 def start_link(cluster) do GenServer.start_link(__MODULE__, cluster) end def init(cluster) do with {:ok, _} <- validate_modules(cluster), {:ok, _} <- validate_behaviours(cluster), {:ok, _} <- activate_tracing(cluster), state <- initialize_monitor(cluster), do: start_monitor(state) end # # Handlers for intercepting :erlang.trace/3 and :erlang.trace_pattern/2 # callbacks for modules registered on the cluster. # def handle_info({:trace_ts, pid, :call, {mod, fun, _args}, ts}, state) do {:noreply, do_intercept_call_request(state, pid, mod, fun, ts)} end # Intercept {:reply, reply, new_state} def handle_info({:trace_ts, pid, :return_from, {_mod, fun, _arity}, {:reply, _, _}, ts}, state) do {:noreply, do_intercept_call_response(state, pid, fun, ts)} end # Intercept {:reply, reply, new_state, timeout | :hibernate} def handle_info({:trace_ts, pid, :return_from, {_mod, fun, _arity}, {:reply, _, _, _}, ts}, state) do {:noreply, do_intercept_call_response(state, pid, fun, ts)} end # Intercept {:noreply, new_state} def handle_info({:trace_ts, pid, :return_from, {_mod, fun, _arity}, {:noreply, _}, ts}, state) do {:noreply, do_intercept_call_response(state, pid, fun, ts)} end # Intercept {:noreply, new_state, timeout | :hibernate} def handle_info({:trace_ts, pid, :return_from, {_mod, fun, _arity}, {:noreply, _, _}, ts}, state) do {:noreply, do_intercept_call_response(state, pid, fun, ts)} end # Intercept {:stop, reason, reply, new_state} def handle_info({:trace_ts, pid, :return_from, {_mod, fun, _arity}, {:stop, _, _, _}, ts}, state) do {:noreply, do_intercept_call_response(state, pid, fun, ts)} end # Intercept {:stop, reason, new_state} def handle_info({:trace_ts, pid, :return_from, {_mod, fun, _arity}, {:stop, _, _}, ts}, state) do {:noreply, do_intercept_call_response(state, pid, fun, ts)} end # Report and rollover metrics window. def handle_info(:rollover_metrics_window, state) do now = :erlang.system_time state = %Monitor{state | duration: Runtime.nano_to_milli(now - state.start)} window = Manager.as_window(state.metrics, statistics?(state)) window = %Window{window | cluster: state.cluster, start: state.start, duration: state.duration} Reporter.push(GenMetrics.GenServer.Reporter, window) Process.send_after(self(), :rollover_metrics_window, window_interval(state)) {:noreply, initialize_monitor(state.cluster, state.metrics)} end # Catch-all for calls not intercepted by monitor. def handle_info(_msg, state), do: {:noreply, state} # # Private utility functions follow. # # Initialize GenServer state for monitor. defp initialize_monitor(cluster, metrics \\ nil) do if metrics do %Monitor{cluster: cluster, metrics: Manager.reinitialize(metrics), start: :erlang.system_time} else %Monitor{cluster: cluster, metrics: Manager.initialize(), start: :erlang.system_time} end end # Initialize periodic callback for metrics reporting and window rollover. defp start_monitor(state) do Process.send_after(self(), :rollover_metrics_window, window_interval(state)) {:ok, state} end # Activate tracing for servers within cluster. defp activate_tracing(cluster) do :erlang.trace(:all, true, [:call, :monotonic_timestamp]) for server <- cluster.servers do :erlang.trace_pattern({server, :handle_call, 3}, [{:_, [], [{:return_trace}]}]) :erlang.trace_pattern({server, :handle_cast, 2}, [{:_, [], [{:return_trace}]}]) :erlang.trace_pattern({server, :handle_info, 2}, [{:_, [], [{:return_trace}]}]) end {:ok, cluster} end # Validate cluster modules can be loaded or report failures. defp validate_modules(cluster) do case require_modules(cluster) do [] -> {:ok, cluster} errs -> {:stop, {:bad_cluster, errs}} end end # Ensure cluster modules are available and can be loaded. defp require_modules(cluster) do cluster.servers |> Enum.uniq |> Runtime.require_modules end # Validate cluster modules implement GenServer or report failures. defp validate_behaviours(cluster) do case require_behaviour(cluster, GenServer) do [] -> {:ok, cluster} errs -> {:stop, {:bad_cluster, errs}} end end # Ensure cluster modules implement GenServer behaviour. defp require_behaviour(cluster, behaviour) do cluster.servers |> Enum.uniq |> Runtime.require_behaviour(behaviour) end defp do_intercept_call_request(state, pid, mod, fun, ts) do if fun in @call_cast_info do do_open_metric(state, pid, mod, fun, ts) else state end end defp do_intercept_call_response(state, pid, fun, ts) do do_close_metric(state, pid, fun, ts) end # Open partial metric on handle_ function call trace. defp do_open_metric(state, pid, mod, fun, ts) do metrics = Manager.open_summary_metric(state.metrics, pid, mod, fun, ts) state = %Monitor{state | metrics: metrics} if statistics?(state) do metrics = Manager.open_stats_metric(state.metrics, pid, mod, fun, ts) %Monitor{state | metrics: metrics} else state end end # Close complete metric on handle_ function return trace. defp do_close_metric(state, pid, events, ts) do metrics = Manager.close_summary_metric(state.metrics, pid, events, ts) state = %Monitor{state | metrics: metrics} if statistics?(state) do metrics = Manager.close_stats_metric(state.metrics, pid, events, ts) %Monitor{state | metrics: metrics} else state end end # Return interval for monitor window rollover. defp window_interval(state) do state.cluster.opts[:window_interval] || @window_interval_default end # Return true if monitor is required to generate optional statistics. defp statistics?(state), do: state.cluster.opts[:statistics] || false end