Current section
8 Versions
Jump to
Current section
8 Versions
Compare versions
5
files changed
+187
additions
-107
deletions
| @@ -5,9 +5,9 @@ | |
| 5 5 | {<<"files">>, |
| 6 6 | [<<"lib">>,<<"lib/orion_collector">>, |
| 7 7 | <<"lib/orion_collector/application.ex">>, |
| 8 | - <<"lib/orion_collector/tracer.ex">>,<<"lib/orion_collector.ex">>, |
| 9 | - <<".formatter.exs">>,<<"mix.exs">>,<<"README.md">>,<<"LICENSE">>, |
| 10 | - <<"CHANGELOG.md">>]}. |
| 8 | + <<"lib/orion_collector/tracer.ex">>,<<"lib/orion_collector/aggregator.ex">>, |
| 9 | + <<"lib/orion_collector.ex">>,<<".formatter.exs">>,<<"mix.exs">>, |
| 10 | + <<"README.md">>,<<"LICENSE">>,<<"CHANGELOG.md">>]}. |
| 11 11 | {<<"licenses">>,[<<"Apache-2.0">>]}. |
| 12 12 | {<<"links">>, |
| 13 13 | [{<<"github">>,<<"https://github.com/LivewareProblems/orion_collector">>}]}. |
| @@ -23,4 +23,4 @@ | |
| 23 23 | {<<"optional">>,false}, |
| 24 24 | {<<"repository">>,<<"hexpm">>}, |
| 25 25 | {<<"requirement">>,<<"~> 1.6">>}]]}. |
| 26 | - {<<"version">>,<<"1.0.1">>}. |
| 26 | + {<<"version">>,<<"1.1.0">>}. |
| @@ -0,0 +1,136 @@ | |
| 1 | + defmodule OrionCollector.Aggregator do |
| 2 | + use GenServer |
| 3 | + import Ex2ms |
| 4 | + |
| 5 | + def start_agg(mfa, pid) do |
| 6 | + spec_args = {OrionCollector.Aggregator, [mfa, pid]} |
| 7 | + DynamicSupervisor.start_child(OrionCollector.AggregatorSupervisor, spec_args) |
| 8 | + end |
| 9 | + |
| 10 | + def mfa_to_name(mfa) do |
| 11 | + {:via, Registry, {OrionCollector.Aggregator.Registry, mfa}} |
| 12 | + end |
| 13 | + |
| 14 | + def stop(pid) do |
| 15 | + GenServer.stop(pid, :normal, 5_000) |
| 16 | + end |
| 17 | + |
| 18 | + def start_link(init = [mfa, pid]) do |
| 19 | + GenServer.start_link(__MODULE__, init, |
| 20 | + name: {:via, Registry, {OrionCollector.Aggregator.Registry, mfa, pid}} |
| 21 | + ) |
| 22 | + end |
| 23 | + |
| 24 | + def child_spec(args) do |
| 25 | + %{ |
| 26 | + id: {OrionCollector.Aggregator, args}, |
| 27 | + start: {OrionCollector.Aggregator, :start_link, [args]}, |
| 28 | + restart: :transient |
| 29 | + } |
| 30 | + end |
| 31 | + |
| 32 | + # --PRIVATE-- |
| 33 | + @impl true |
| 34 | + def init([mfa, pid]) do |
| 35 | + mon_ref = Process.monitor(pid) |
| 36 | + |
| 37 | + :erlang.trace_pattern(mfa, match_spec(), [:local]) |
| 38 | + |
| 39 | + Process.send_after(self(), :send_data, 500) |
| 40 | + |
| 41 | + initial_state = %{ |
| 42 | + mfa: mfa, |
| 43 | + liveview_pid: pid, |
| 44 | + call_depth: %{}, |
| 45 | + time_stored: %{}, |
| 46 | + ddsketch: DogSketch.SimpleDog.new(), |
| 47 | + ref_mon: mon_ref |
| 48 | + } |
| 49 | + |
| 50 | + {:ok, initial_state} |
| 51 | + end |
| 52 | + |
| 53 | + @impl true |
| 54 | + def handle_cast( |
| 55 | + {:trace_ts, trace_pid, :call, _mfa, start_time}, |
| 56 | + state = %{call_depth: cd_map, time_stored: time_stored_map} |
| 57 | + ) do |
| 58 | + cd = Map.get(cd_map, trace_pid, 0) |
| 59 | + new_cd_map = Map.put(cd_map, trace_pid, cd + 1) |
| 60 | + |
| 61 | + new_ts_map = |
| 62 | + if cd == 0 do |
| 63 | + Map.put(time_stored_map, {trace_pid, cd + 1}, start_time) |
| 64 | + else |
| 65 | + time_stored_map |
| 66 | + end |
| 67 | + |
| 68 | + new_state = |
| 69 | + state |
| 70 | + |> Map.put(:call_depth, new_cd_map) |
| 71 | + |> Map.put(:time_stored, new_ts_map) |
| 72 | + |
| 73 | + {:noreply, new_state} |
| 74 | + end |
| 75 | + |
| 76 | + @accepted_return_tags [:return_from, :exception_from] |
| 77 | + |
| 78 | + @impl true |
| 79 | + def handle_cast( |
| 80 | + {:trace_ts, trace_pid, return_tag, _mfa, _TraceTerm, end_time}, |
| 81 | + state = %{call_depth: cd_map, time_stored: time_stored_map, ddsketch: ddsketch} |
| 82 | + ) |
| 83 | + when return_tag in @accepted_return_tags do |
| 84 | + case Map.get(cd_map, trace_pid, 0) do |
| 85 | + 0 -> |
| 86 | + {:noreply, state} |
| 87 | + |
| 88 | + 1 -> |
| 89 | + new_cd_map = Map.delete(cd_map, trace_pid) |
| 90 | + {start_time, new_ts_map} = Map.pop(time_stored_map, {trace_pid, 1}) |
| 91 | + |
| 92 | + call_time_micro = :timer.now_diff(end_time, start_time) |
| 93 | + new_sketch = DogSketch.SimpleDog.insert(ddsketch, call_time_micro / 1_000) |
| 94 | + |
| 95 | + new_state = |
| 96 | + state |
| 97 | + |> Map.put(:call_depth, new_cd_map) |
| 98 | + |> Map.put(:time_stored, new_ts_map) |
| 99 | + |> Map.put(:ddsketch, new_sketch) |
| 100 | + |
| 101 | + {:noreply, new_state} |
| 102 | + |
| 103 | + cd when cd > 1 -> |
| 104 | + new_cd_map = Map.put(cd_map, trace_pid, cd - 1) |
| 105 | + |
| 106 | + {:noreply, Map.put(state, :call_depth, new_cd_map)} |
| 107 | + end |
| 108 | + end |
| 109 | + |
| 110 | + @impl true |
| 111 | + def handle_info(:send_data, %{ddsketch: ddsketch, liveview_pid: liveview_pid} = state) do |
| 112 | + new_sketch = DogSketch.SimpleDog.new() |
| 113 | + |
| 114 | + if ddsketch != new_sketch do |
| 115 | + send(liveview_pid, {:ddsketch, ddsketch}) |
| 116 | + end |
| 117 | + |
| 118 | + Process.send_after(self(), :send_data, 500) |
| 119 | + |
| 120 | + {:noreply, Map.put(state, :ddsketch, new_sketch)} |
| 121 | + end |
| 122 | + |
| 123 | + @impl true |
| 124 | + def handle_info({:DOWN, ref, :process, _object, _reason}, %{ref_mon: ref, mfa: mfa} = state) do |
| 125 | + :erlang.trace_pattern(mfa, false, []) |
| 126 | + {:stop, :normal, state} |
| 127 | + end |
| 128 | + |
| 129 | + defp match_spec() do |
| 130 | + fun do |
| 131 | + _ -> |
| 132 | + return_trace() |
| 133 | + exception_trace() |
| 134 | + end |
| 135 | + end |
| 136 | + end |
| @@ -8,7 +8,9 @@ defmodule OrionCollector.Application do | |
| 8 8 | @impl true |
| 9 9 | def start(_type, _args) do |
| 10 10 | children = [ |
| 11 | - {DynamicSupervisor, strategy: :one_for_one, name: OrionCollector.TracerSupervisor} |
| 11 | + {Registry, keys: :unique, name: OrionCollector.Aggregator.Registry}, |
| 12 | + OrionCollector.Tracer, |
| 13 | + {DynamicSupervisor, strategy: :one_for_one, name: OrionCollector.AggregatorSupervisor} |
| 12 14 | |
| 13 15 | # Starts a worker by calling: OrionCollector.Worker.start_link(arg) |
| 14 16 | # {OrionCollector.Worker, arg} |
| @@ -16,7 +18,7 @@ defmodule OrionCollector.Application do | |
| 16 18 | |
| 17 19 | # See https://hexdocs.pm/elixir/Supervisor.html |
| 18 20 | # for other strategies and supported options |
| 19 | - opts = [strategy: :one_for_one, name: OrionCollector.Supervisor] |
| 21 | + opts = [strategy: :rest_for_one, name: OrionCollector.Supervisor] |
| 20 22 | Supervisor.start_link(children, opts) |
| 21 23 | end |
| 22 24 | end |
| @@ -1,7 +1,13 @@ | |
| 1 1 | defmodule OrionCollector.Tracer do |
| 2 2 | use GenServer |
| 3 | - import Ex2ms |
| 4 3 | |
| 4 | + alias OrionCollector.Aggregator |
| 5 | + |
| 6 | + @moduledoc """ |
| 7 | + This is the process that collect the Trace, then dispatch them per aggregator |
| 8 | + |
| 9 | + There is one Tracer per node, but one Aggregator per MFA being traced. |
| 10 | + """ |
| 5 11 | def start_all_node_tracers(mfa, self, start_status \\ :running) do |
| 6 12 | if self do |
| 7 13 | OrionCollector.Tracer.start_tracer(mfa, self(), start_status) |
| @@ -22,31 +28,30 @@ defmodule OrionCollector.Tracer do | |
| 22 28 | :ok |
| 23 29 | end |
| 24 30 | |
| 25 | - def pause_trace(mfa, self) do |
| 31 | + def pause_trace(self) do |
| 26 32 | if self do |
| 27 | - :erlang.trace_pattern(mfa, false, [:local]) |
| 33 | + OrionCollector.Tracer.change_status(:pause) |
| 28 34 | end |
| 29 35 | |
| 30 | - :erpc.multicall(list_nodes(), :erlang, :trace_pattern, [mfa, false, [:local]], 5_000) |
| 36 | + :erpc.multicall(list_nodes(), OrionCollector.Tracer, :change_status, [:pause], 5_000) |
| 31 37 | end |
| 32 38 | |
| 33 | - def restart_trace(mfa, self) do |
| 39 | + def restart_trace(self) do |
| 34 40 | if self do |
| 35 | - :erlang.trace_pattern(mfa, match_spec(), [:local]) |
| 41 | + OrionCollector.Tracer.change_status(:start) |
| 36 42 | end |
| 37 43 | |
| 38 | - :erpc.multicall( |
| 39 | - list_nodes(), |
| 40 | - :erlang, |
| 41 | - :trace_pattern, |
| 42 | - [mfa, match_spec(), [:local]], |
| 43 | - 5_000 |
| 44 | - ) |
| 44 | + :erpc.multicall(list_nodes(), OrionCollector.Tracer, :change_status, [:start], 5_000) |
| 45 45 | end |
| 46 46 | |
| 47 47 | def start_tracer(mfa, pid, start_status) do |
| 48 | - spec_args = {OrionCollector.Tracer, [mfa, pid, start_status]} |
| 49 | - DynamicSupervisor.start_child(OrionCollector.TracerSupervisor, spec_args) |
| 48 | + Aggregator.start_agg(mfa, pid) |
| 49 | + change_status(start_status) |
| 50 | + end |
| 51 | + |
| 52 | + @spec change_status(any) :: any |
| 53 | + def change_status(status) do |
| 54 | + GenServer.call(__MODULE__, status) |
| 50 55 | end |
| 51 56 | |
| 52 57 | def stop(pid) do |
| @@ -54,119 +59,56 @@ defmodule OrionCollector.Tracer do | |
| 54 59 | end |
| 55 60 | |
| 56 61 | def start_link(init \\ []) do |
| 57 | - GenServer.start_link(__MODULE__, init) |
| 62 | + GenServer.start_link(__MODULE__, init, name: __MODULE__) |
| 58 63 | end |
| 59 64 | |
| 60 65 | def child_spec(args) do |
| 61 66 | %{ |
| 62 | - id: {OrionCollector.Tracer, args}, |
| 67 | + id: OrionCollector.Tracer, |
| 63 68 | start: {OrionCollector.Tracer, :start_link, [args]}, |
| 64 | - restart: :transient |
| 69 | + restart: :permanent |
| 65 70 | } |
| 66 71 | end |
| 67 72 | |
| 68 73 | # --PRIVATE-- |
| 69 74 | @impl true |
| 70 | - def init([mfa, pid, start_status]) do |
| 71 | - mon_ref = Process.monitor(pid) |
| 75 | + def init(_start) do |
| 76 | + running_trace(false) |
| 72 77 | |
| 73 | - if start_status == :running do |
| 74 | - :erlang.trace_pattern(mfa, match_spec(), [:local]) |
| 75 | - end |
| 76 | - |
| 77 | - :erlang.trace(:all, true, [:call, :arity, :timestamp]) |
| 78 | - |
| 79 | - Process.send_after(self(), :send_data, 500) |
| 80 | - |
| 81 | - initial_state = %{ |
| 82 | - mfa: mfa, |
| 83 | - liveview_pid: pid, |
| 84 | - call_depth: %{}, |
| 85 | - time_stored: %{}, |
| 86 | - ddsketch: DogSketch.SimpleDog.new(), |
| 87 | - ref_mon: mon_ref |
| 88 | - } |
| 78 | + initial_state = %{running_status: :pause} |
| 89 79 | |
| 90 80 | {:ok, initial_state} |
| 91 81 | end |
| 92 82 | |
| 83 | + @impl true |
| 84 | + def handle_call(start_status, _from, state) do |
| 85 | + running_trace(start_status == :start) |
| 86 | + {:reply, :ok, Map.put(state, :running_status, start_status)} |
| 87 | + end |
| 88 | + |
| 93 89 | @impl true |
| 94 90 | def handle_info( |
| 95 | - {:trace_ts, trace_pid, :call, _mfa, start_time}, |
| 96 | - state = %{call_depth: cd_map, time_stored: time_stored_map} |
| 91 | + {:trace_ts, _trace_pid, :call, mfa, _start_time} = trace_msg, |
| 92 | + state |
| 97 93 | ) do |
| 98 | - cd = Map.get(cd_map, trace_pid, 0) |
| 99 | - new_cd_map = Map.put(cd_map, trace_pid, cd + 1) |
| 100 | - |
| 101 | - new_ts_map = |
| 102 | - if cd == 0 do |
| 103 | - Map.put(time_stored_map, {trace_pid, cd + 1}, start_time) |
| 104 | - else |
| 105 | - time_stored_map |
| 106 | - end |
| 107 | - |
| 108 | - new_state = |
| 109 | - state |
| 110 | - |> Map.put(:call_depth, new_cd_map) |
| 111 | - |> Map.put(:time_stored, new_ts_map) |
| 112 | - |
| 113 | - {:noreply, new_state} |
| 94 | + GenServer.cast(Aggregator.mfa_to_name(mfa), trace_msg) |
| 95 | + {:noreply, state} |
| 114 96 | end |
| 115 97 | |
| 116 98 | @accepted_return_tags [:return_from, :exception_from] |
| 117 99 | |
| 118 100 | @impl true |
| 119 101 | def handle_info( |
| 120 | - {:trace_ts, trace_pid, return_tag, _mfa, _TraceTerm, end_time}, |
| 121 | - state = %{call_depth: cd_map, time_stored: time_stored_map, ddsketch: ddsketch} |
| 102 | + {:trace_ts, _trace_pid, return_tag, mfa, _TraceTerm, _end_time} = trace_msg, |
| 103 | + state |
| 122 104 | ) |
| 123 105 | when return_tag in @accepted_return_tags do |
| 124 | - case Map.get(cd_map, trace_pid, 0) do |
| 125 | - 0 -> |
| 126 | - {:noreply, state} |
| 127 | - |
| 128 | - 1 -> |
| 129 | - new_cd_map = Map.delete(cd_map, trace_pid) |
| 130 | - {start_time, new_ts_map} = Map.pop(time_stored_map, {trace_pid, 1}) |
| 131 | - |
| 132 | - call_time_micro = :timer.now_diff(end_time, start_time) |
| 133 | - new_sketch = DogSketch.SimpleDog.insert(ddsketch, call_time_micro / 1_000) |
| 134 | - |
| 135 | - new_state = |
| 136 | - state |
| 137 | - |> Map.put(:call_depth, new_cd_map) |
| 138 | - |> Map.put(:time_stored, new_ts_map) |
| 139 | - |> Map.put(:ddsketch, new_sketch) |
| 140 | - |
| 141 | - {:noreply, new_state} |
| 142 | - |
| 143 | - cd when cd > 1 -> |
| 144 | - new_cd_map = Map.put(cd_map, trace_pid, cd - 1) |
| 145 | - |
| 146 | - {:noreply, Map.put(state, :call_depth, new_cd_map)} |
| 147 | - end |
| 106 | + GenServer.cast(Aggregator.mfa_to_name(mfa), trace_msg) |
| 107 | + {:noreply, state} |
| 148 108 | end |
| 149 109 | |
| 150 | - @impl true |
| 151 | - def handle_info(:send_data, %{ddsketch: ddsketch, liveview_pid: liveview_pid} = state) do |
| 152 | - send(liveview_pid, {:ddsketch, ddsketch}) |
| 153 | - Process.send_after(self(), :send_data, 500) |
| 154 | - |
| 155 | - {:noreply, Map.put(state, :ddsketch, DogSketch.SimpleDog.new())} |
| 156 | - end |
| 157 | - |
| 158 | - @impl true |
| 159 | - def handle_info({:DOWN, ref, :process, _object, _reason}, %{ref_mon: ref, mfa: mfa} = state) do |
| 160 | - :erlang.trace_pattern(mfa, false, []) |
| 161 | - {:stop, :normal, state} |
| 162 | - end |
| 163 | - |
| 164 | - defp match_spec() do |
| 165 | - fun do |
| 166 | - _ -> |
| 167 | - return_trace() |
| 168 | - exception_trace() |
| 169 | - end |
| 110 | + defp running_trace(bool) do |
| 111 | + :erlang.trace(:all, bool, [:call, :arity, :timestamp]) |
| 170 112 | end |
| 171 113 | |
| 172 114 | defp list_nodes() do |
| @@ -1,7 +1,7 @@ | |
| 1 1 | defmodule OrionCollector.MixProject do |
| 2 2 | use Mix.Project |
| 3 3 | |
| 4 | - @version "1.0.1" |
| 4 | + @version "1.1.0" |
| 5 5 | |
| 6 6 | def project do |
| 7 7 | [ |